为何修改Apache Beam的ParDo中yield为return后输出不同?
Apache Beam ParDo中yield与return的输出差异问题
我正在使用Python SDK学习Apache Beam,在查阅官网代码后产生疑问:
官网原代码
import apache_beam as beam import re class SplitWords(beam.DoFn): def __init__(self, delimiter=','): self.delimiter = delimiter def process(self, text): for word in text.split(self.delimiter): yield word with beam.Pipeline() as pipeline: plants = ( pipeline | 'Gardening plants' >> beam.Create([ '🍓Strawberry,🥕Carrot,🍆Eggplant', '🍅Tomato,🥔Potato', ]) | 'Split words' >> beam.ParDo(SplitWords(',')) | beam.Map(print))
原代码输出
🍓Strawberry 🥕Carrot 🍆Eggplant 🍅Tomato 🥔Potato
修改后的代码(yield替换为return)
import apache_beam as beam import re class SplitWords(beam.DoFn): def __init__(self, delimiter=','): self.delimiter = delimiter def process(self, text): for word in text.split(self.delimiter): return [word] with beam.Pipeline() as pipeline: plants = ( pipeline | 'Gardening plants' >> beam.Create([ '🍓Strawberry,🥕Carrot,🍆Eggplant', '🍅Tomato,🥔Potato', ]) | 'Split words' >> beam.ParDo(SplitWords(',')) | beam.Map(print))
修改后输出
🍓Strawberry 🍅Tomato
请问为何修改后输出结果不同?
解答
问题出在你修改后的process方法逻辑上:
- 原代码用
yield时,for循环会遍历拆分后的每个单词,每次yield都会向Pipeline输出一个元素,因此每个输入文本的所有拆分结果都会被输出。 - 你修改后的代码在
for循环内部使用return [word],第一次循环迭代就会触发return,直接终止函数执行,后续的循环步骤根本不会运行。这就导致每个输入字符串只返回第一个拆分出的元素,最终输出只有两个结果。
Beam文档说的"return语句返回一个可迭代对象",是指你要返回包含所有结果的可迭代对象,而不是在循环里逐个return。正确的写法应该直接返回拆分后的整个列表:
def process(self, text): return text.split(self.delimiter)
或者返回生成器表达式也可以:
def process(self, text): return (word for word in text.split(self.delimiter))
这样修改后,每个输入文本的所有拆分结果都会被输出,和原代码的yield效果一致。
内容的提问来源于stack exchange,提问作者Amarjeet Kushwaha
相关产品推荐
相关产品推荐

