Apache Beam:JSON列表转多元素PCollection遇仅返回首元素问题求解
解决Apache Beam中JSON列表仅解析第一个元素的问题
这个问题的核心是转换操作的选择和函数返回逻辑出了问题,咱们一步步来修正:
问题根源
你现在用的beam.Map是一对一的转换:每个输入文本行对应一个输出元素。而你的parse_json函数在循环里第一次迭代就用return i直接返回了,导致后续的列表元素根本没机会被处理,自然只能拿到第一个元素。
要实现一个输入(JSON字符串)对应多个输出(列表中的每个元素),咱们需要用beam.FlatMap,它专门处理这种一对多的展开场景。
修改方案
有两种简洁的实现方式,选哪种都可以:
方式一:用生成器逐个产出元素
把函数里的return改成yield,让函数变成一个生成器,逐个返回列表里的每个元素:
def parse_json(data): import json for i in json.loads(data): yield i
方式二:直接返回解析后的列表
因为FlatMap会自动迭代可迭代对象(比如列表)并输出每个元素,所以可以简化函数:
def parse_json(data): import json return json.loads(data)
最后替换转换操作
把原来的beam.Map换成beam.FlatMap:
data = (p | "Read text" >> beam.io.textio.ReadFromText(f'gs://{bucket_name}/not_processed/2020-06-08T23:59:59.999Z__rms004_m1__not_sent_msg.txt') | "Parse json" >> beam.FlatMap(parse_json))
为什么这样能行?
FlatMap的工作逻辑是:对每个输入元素,调用你的函数后,如果函数返回了一个可迭代对象(生成器、列表等),它会把这个对象里的每一个元素单独提取出来,作为PCollection的独立元素。这样原来的JSON列表里的所有元素就都能被正确展开到你的PCollection里了。
内容的提问来源于stack exchange,提问作者mmc
相关产品推荐
相关产品推荐

