DataflowRunner下Apache Beam出现AttributeError: tuple无encode属性排查
问题:DataflowRunner下有状态DoFn的编码错误
我在使用有状态且基于时间的DoFn,希望在固定窗口结束后2秒处理数据。已在Apache Beam Playground用DirectRunner测试过可复现的代码示例,输入DoFn的数据格式为KV[str, str],本地代码使用DataflowRunner,这是唯一的环境差异。
转换格式的DoFn(AddKeys)
class AddKeys(beam.DoFn): def __init__(self, settings): self.settings = settings def process(self, element): data = element["data"] for setting in self.settings: if setting["data"] == data: yield [ ("tuple key 1", setting["tuple key 1"]), ("tuple key 2", setting["tuple key 2"]), ("tuple key 3", setting["tuple key 3"]), ("element", str(element)) ]
有状态DoFn(ProcessCollection)
class ProcessCollection(beam.DoFn): EXPIRY_TIMER = TimerSpec('expiry', TimeDomain.WATERMARK) BUFFER_STATE = BagStateSpec( 'buffer', ListCoder(StrUtf8Coder())) def process(self, element, timer=beam.DoFn.TimerParam(EXPIRY_TIMER), window=beam.DoFn.WindowParam, buffer=beam.DoFn.StateParam(BUFFER_STATE)): timer.set(window.end + Duration(seconds=2)) buffer.add(str(element)) @on_timer(EXPIRY_TIMER) def expiry(self, buffer=beam.DoFn.StateParam(BUFFER_STATE)): events = buffer.read() for event in events: yield ''.join(event) buffer.clear()
Pipeline调用代码
# 后续处理并分支到有状态DoFn extra_processing = ( raw_data_processing | "Add Group Keys" >> beam.Map( lambda message: add_group_key( message, SETTINGS) ) | "Fixed Window" >> beam.WindowInto( window.FixedWindows(self.window_length), # 允许700ms延迟的数据 allowed_lateness=window.Duration(seconds=0.7) ) | "Group" >> beam.GroupByKey() | "Process Further" >> beam.ParDo(OtherDoFn(SETTINGS, CONFIG)) ) # 使用有状态DoFn处理数据 ( extra_processing | "Add Keys" >> beam.ParDo(AddKeys(SETTINGS)).with_output_types(KV[str, str]) | "Process Collection" >> beam.ParDo(ProcessCollection()) | 'Log' >> beam.LogElements(with_timestamp=True) )
错误信息
File "/usr/local/lib/python3.10/site-packages/apache_beam/coders/coders.py", line 429, in encode
return value.encode('utf-8')
AttributeError: 'tuple' object has no attribute 'encode' [while running 'Add Keys-ptransform-51']
问题原因与修复方案
核心原因
错误源于AddKeys DoFn的输出与声明的类型不匹配:
- 你通过
.with_output_types(KV[str, str])声明输出是键值对,但AddKeys的process方法中yield的是包含多个元组的列表,而非单个KV元素。 - 后续编码时,整个列表被当作元素传递,列表中的元组被误判为需要编码的字符串对象,尝试调用
encode方法时触发错误(元组没有该方法)。 - DirectRunner与DataflowRunner的类型校验严格度不同,Playground的DirectRunner未严格校验,因此未触发错误。
修复方案
修改AddKeys的process方法,将列表中的每个元组单独yield,确保每个输出元素都是符合KV[str, str]的单个键值对:
class AddKeys(beam.DoFn): def __init__(self, settings): self.settings = settings def process(self, element): data = element["data"] for setting in self.settings: if setting["data"] == data: # 逐个输出元组,每个元组对应一个KV元素 yield ("tuple key 1", setting["tuple key 1"]) yield ("tuple key 2", setting["tuple key 2"]) yield ("tuple key 3", setting["tuple key 3"]) yield ("element", str(element))
这样每个输出元素都是(str, str)类型的元组,符合声明的输出类型,编码时就能正确处理字符串的encode操作。
内容的提问来源于stack exchange,提问作者Joe Moore
相关产品推荐
相关产品推荐

