在Apache NiFi中如何将两个CSV合并为JSON并合并模式
使用Apache NiFi合并两个CSV为嵌套JSON格式
我来帮你梳理下具体实现步骤,结合你现有的GetFile → UpdateAttribute → Split Records流程来扩展调整,最终实现将人员信息和关联事件合并成带事件数组的JSON结构:
1. 先准备Schema文件
NiFi的Record组件依赖Avro Schema来定义数据结构,你需要创建三个Schema文件:
persons_schema.avsc(人员数据结构)
{ "type": "record", "name": "Person", "fields": [ {"name": "Id", "type": "string"}, {"name": "Name", "type": "string"}, {"name": "Surname", "type": "string"} ] }
events_schema.avsc(事件数据结构)
{ "type": "record", "name": "Event", "fields": [ {"name": "EId", "type": "int"}, {"name": "PId", "type": "string"}, {"name": "Date", "type": "string"}, {"name": "Desc", "type": "string"} ] }
merged_person_events_schema.avsc(合并后的嵌套结构)
{ "type": "record", "name": "PersonWithEvents", "fields": [ {"name": "Id", "type": "string"}, {"name": "Name", "type": "string"}, {"name": "Surname", "type": "string"}, {"name": "Events", "type": {"type": "array", "items": "Event"}} ] }
2. 完整NiFi流程构建
我分两条分支处理两个CSV,最后关联合并:
分支一:处理人员数据并缓存
- GetFile:配置读取
persons.csv,指定文件所在的输入目录。 - UpdateAttribute:添加两个关键属性:
schema.name:值为你存放persons_schema.avsc的绝对路径(比如/nifi/schemas/persons_schema.avsc)field.delimiter:设为|(因为你的CSV用竖线分隔)
- SplitRecords:选择
CSVRecordReader,勾选First Line is Header,引用上面的人员Schema,把整个CSV拆分成单条人员记录。 - PutDistributedMapCache:先启用
EmbeddedDistributedMapCacheServer服务,然后配置这个组件:Cache Key:填${record:value('/Id')}(用人员ID作为缓存key)Cache Value:填${record:toJson()}(把单条人员记录转成JSON存在缓存里)
分支二:处理事件数据并关联人员
- GetFile:配置读取
events.csv,指定单独的输入目录。 - UpdateAttribute:添加属性:
schema.name:值为events_schema.avsc的绝对路径field.delimiter:设为|
- SplitRecords:同样用
CSVRecordReader,勾选表头选项,拆分出单条事件记录。 - LookupRecord:
- 选择
DistributedMapCacheRecordLookupService作为查询服务 Lookup Key:填${record:value('/PId')}(用事件里的人员ID匹配缓存中的key)- 返回的Record Schema选择
persons_schema.avsc,这样每条事件记录会带上对应的人员信息
- 选择
- MergeRecord:把同一人员的所有事件合并成数组:
Merge Strategy选Merge Records by Specific Fields,Merge Fields填Id(人员ID)Record Writer选JsonRecordSetWriter,引用merged_person_events_schema.avsc,这样会自动把多条事件合并到Events数组中
- PutFile:把最终的JSON文件输出到你指定的目录即可
3. 替代方案(无需缓存)
如果你的文件数据量不大,可以不用缓存,直接用JoinRecord组件替代缓存+查询的步骤:
- 把两个
SplitRecords输出的人员和事件记录都发送到JoinRecord - 设置
Join Type为Left Outer Join,关联字段为人员的Id和事件的PId - 之后再用
MergeRecord合并同一人员的事件数组,流程更简洁
最终输出示例
合并后的JSON格式如下:
[ { "Id": "ABC-123", "Name": "John", "Surname": "Smith", "Events": [ { "EId": 1, "PId": "ABC-123", "Date": "2017-05-01", "Desc": "Groove party" }, { "EId": 3, "PId": "ABC-123", "Date": "2017-06-01", "Desc": "The night out" } ] }, { "Id": "ABC-111", "Name": "Allan", "Surname": "Wood", "Events": [ { "EId": 2, "PId": "ABC-111", "Date": "2017-06-01", "Desc": "Snack No. One" } ] }, { "Id": "ABC-001", "Name": "Grace", "Surname": "Kelly", "Events": [] } ]
内容的提问来源于stack exchange,提问作者Vaasha
相关产品推荐
相关产品推荐

