【问题标题】:Apache Nifi MergeContent output data inconsistent?Apache Nifi MergeContent 输出数据不一致?
【发布时间】:2018-06-05 14:28:45
【问题描述】:

对使用 nifi 还是很陌生。在设计方面需要帮助。 我正在尝试在 HDFS 目录中使用虚拟 csv 文件(目前)创建一个简单的流程,并将一些文本数据添加到每个流程文件中的每个记录中。

传入文件:

dummy1.csv
dummy2.csv
dummy3.csv

内容:

"Eldon Base for stackable storage shelf, platinum",Muhammed MacIntyre,3,-213.25,38.94,35,Nunavut,Storage & Organization,0.8
"1.7 Cubic Foot Compact ""Cube"" Office Refrigerators",BarryFrench,293,457.81,208.16,68.02,Nunavut,Appliances,0.58
"Cardinal Slant-D Ring Binder, Heavy Gauge Vinyl",Barry French,293,46.71,8.69,2.99,Nunavut,Binders and Binder Accessories,0.39
...

期望的输出:

d17a3259-0718-4c7b-bee8-924266aebcc7,Mon Jun 04 16:36:56 EDT 2018,Fellowes Recycled Storage Drawers,Allen Rosenblatt,11137,395.12,111.03,8.64,Northwest Territories,Storage & Organization,0.78
25f17667-9216-4f1d-b69c-23403cd13464,Mon Jun 04 16:36:56 EDT 2018,Satellite Sectional Post Binders,Barry Weirich,11202,79.59,43.41,2.99,Northwest Territories,Binders and Binder Accessories,0.39
ce0b569f-5d93-4a54-b55e-09c18705f973,Mon Jun 04 16:36:56 EDT 2018,Deflect-o DuraMat Antistatic Studded Beveled Mat for Medium Pile Carpeting,Doug Bickford,11456,399.37,105.34,24.49,Northwest Territories,Office Furnishings,0.61

流程 splitText- 替换文本- 合并内容-

(这可能是实现我想要获得的效果的一种糟糕方法,但我在某处看到 uuid 是生成唯一会话 id 的最佳选择。因此考虑将每一行从传入数据提取到流文件并生成uuid)

但是不知何故,您可以看到数据的顺序混乱了。前 3 行的输出不同。但是,我正在使用的测试数据(50000 个条目)似乎在其他行中有数据。多次测试通常显示 2001 行之后的数据顺序发生变化。

是的,我确实在这里搜索了类似的问题,并尝试在合并中使用碎片整理方法,但它没有用。如果有人可以解释这里发生的事情以及如何使用唯一的 session_id,每个记录的时间戳以相同的方式获取数据,我将不胜感激。我需要更改或修改某些参数以获得正确的输出吗?如果有更好的方法,我也愿意接受建议。

【问题讨论】:

    标签: hadoop hdfs cloudera apache-nifi hortonworks-data-platform


    【解决方案1】:

    首先感谢您如此详尽而详细的回复。我想你解决了我对处理器工作原理的很多疑问!

    只有在碎片整理模式下才能保证合并的顺序,因为它将根据它们的片段索引将流文件按顺序排列。我不确定为什么这不起作用,但如果您可以使用显示问题的示例数据创建一个流程模板,这将有助于调试。

    我将尝试再次使用干净的模板复制此方法。可能是一些参数问题,HDFS 写入器无法写入。

    我不确定您的流程的目的是重新合并拆分的原始 CSV,还是将几个不同的 CSV 合并在一起。碎片整理模式只会重新合并原始 CSV,因此如果 ListHDFS 拾取 10 个 CSV,在拆分和重新合并后,您应该再次拥有 10 个 CSV。

    是的,这正是我所需要的。将数据拆分并连接到其相应的文件。我还没有特别需要再次加入输出。

    将 CSV 拆分为每个流文件 1 行以操作每一行的方法是一种常见的方法,但是如果您有许多大型 CSV 文件,它就不会很好地执行。一种更有效的方法是尝试在原地操作数据而不进行拆分。这通常可以使用面向记录的处理器来完成。

    1. 我纯粹是本能地使用这种方法,并没有意识到这是一种常见的方法。有时数据文件可能非常大,这意味着单个文件中有超过一百万条记录。这不会是集群中的 i/o 问题吗?因为这意味着每条记录=一个流文件=一个唯一的 uuid。 nifi 可以处理多少个合适的流文件? (我知道这取决于集群配置,并会尝试从 hdp 管理员那里获取有关集群的更多信息)
    2. 您对“尝试在原地操作数据而不进行拆分”有何建议?你能给出一个例子或模板或处理器来使用吗?

    在这种情况下,您需要为 CSV 定义一个架构,其中包括数据中的所有列,以及会话 ID 和时间戳。然后使用 UpdateRecord 处理器,您将使用像 /session_id = ${UUID()} 和 /timestamp = ${now()} 这样的记录路径表达式。这将逐行流式传输内容并更新每条记录并将其写回,将其全部保存为一个流文件。

    这看起来很有希望。你能分享一个简单的模板从hdfs中提取文件>处理>写入hdfs文件但不分割吗?

    由于限制,我不愿分享模板。但是让我看看我是否可以创建一个通用的模板,我会分享

    感谢您的智慧! :)

    【讨论】:

    • @StrangerThinks 这里是一个模板,展示了如何在不拆分的情况下更新 CSV 记录 - gist.githubusercontent.com/bbende/…
    • 如果您在使用${now()} 时遇到问题,请尝试${now():format('yyyy-MM-dd HH:mm:ss')},我在类型之间转换时遇到了一些问题,这有帮助
    • 模板定义了一个本地 AvroSchemaRegistry,这是模式所在的位置,您可以使用左侧的上下文面板从您所在的进程组的控制器服务部分访问它。跨度>
    • 我使用了我昨天创建的模板,并根据您上面的内容更新了架构,然后我更改了 GenerateFlowFile 处理器自定义文本,以使您在原始问题中添加 3 行 CSV 行,并且它正确地将会话和时间戳添加到每一行,并在输出中包含完整的行
    • 很难说没有您的实际数据...如果您使用我模板中的模式,那么它将读取每行的前两列,在您的数据中是产品和名字,然后它将添加会话和时间戳,并应写出 4 列(会话、时间戳、产品、名字)。不知道为什么完整架构对您不起作用,除非您的数据不符合它。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-11-17
    • 1970-01-01
    • 2016-11-13
    • 1970-01-01
    相关资源
    最近更新 更多