求助:使用Apache Beam去除JSON文件重复数据失败
问题分析与解决方案
你的代码没法实现去重的核心原因是函数内部的临时变量(temp、res)每次处理元素都会重置,根本没法跨元素记录已经出现过的idtrx,而且后续的判断逻辑也存在错误(比如element[key]的key是循环最后一个键,不一定是idtrx)。
Apache Beam是分布式处理框架,去重需要用框架提供的分布式变换,而不是单元素的处理函数。下面给你两种可行的实现方式:
方法1:使用Beam内置的Distinct变换(推荐)
Distinct可以直接对PCollection去重,通过指定key_fn来基于idtrx判断重复项:
import apache_beam as beam with beam.Pipeline() as p: ( p | "Read JSON" >> beam.io.ReadFromText("input.json") | "Parse JSON" >> beam.Map(lambda x: eval(x.strip().rstrip(','))) # 处理输入的逗号分隔格式 | "Deduplicate by idtrx" >> beam.Distinct(lambda element: element['idtrx']) | "Write Result" >> beam.io.WriteToText("output.json") )
注:如果你的JSON输入是每行一个对象,Parse JSON可以换成beam.Map(json.loads);如果是你提供的逗号分隔格式,需要先处理掉每个对象末尾的逗号再解析。
方法2:使用GroupByKey+取首元素
如果需要保留重复项中的某一个(比如第一个出现的完整对象),可以先按idtrx分组,再取每组的第一个元素:
import apache_beam as beam class ExtractIdtrx(beam.DoFn): def process(self, element): yield (element['idtrx'], element) with beam.Pipeline() as p: ( p | "Read JSON" >> beam.io.ReadFromText("input.json") | "Parse JSON" >> beam.Map(lambda x: eval(x.strip().rstrip(','))) | "Extract idtrx as key" >> beam.ParDo(ExtractIdtrx()) | "Group by idtrx" >> beam.GroupByKey() | "Take first element in group" >> beam.Map(lambda x: x[1][0]) | "Write Result" >> beam.io.WriteToText("output.json") )
原代码的错误点拆解
temp和res是函数局部变量,每次调用CleansingDuplicateData都会重新创建,无法记住之前处理过的idtrx。- 第二个
if块中的element[key],key是上一次循环的最后一个键(比如UserCreatedDate),不是你要判断的idtrx,逻辑完全错误。 - 最后返回原
element,没有任何过滤重复的操作,等于没处理。
内容的提问来源于stack exchange,提问作者Rafif Dwiputra
相关产品推荐
相关产品推荐

