Код: Выделить всё
# main.py
with Pipeline(options=options) as pipeline:
_ = (
pipeline
| "Read input file" >> ReadFromText(options.input_file)
| "Execute" >> Execute()
| "Write CSV File" >> WriteToText(options.output_file)
)
# execute.py
class Execute(beam.PTransform):
def expand(self, pcollection: PCollection[str]) -> PCollection[str]:
return pcollection | "Execute" >> beam.ParDo(ExecuteFn())
class ExecuteFn(beam.DoFn):
def __init__(self):
self._lines: list[str] = []
def process(self, element: str) -> Iterator[str]:
self._lines.append(element)
yield element
def start_bundle(self) -> None:
log.error("Starting bundle")
self._lines = []
def finish_bundle(self) -> None:
log.error(f"Finishing bundle, outputting {len(self._lines)} lines")
for line in self._lines:
log.error(f"Outputting line: {line}")
Ожидался запуск кода в прямом бегуне и в модульных тестах, но в потоке данных start_bundle и Finish_bundle никогда не вызываются.
Что я пробовал:
- Чтобы гарантировать, что Finish_bundle действительно никогда не вызывается, я добавил операторы журнала в start_bundle и Finish_bundle. В модульных тестах и при использовании прямого выполнения эти сообщения отображаются по мере выполнения поведения. В потоке данных GCP ничего не происходит. Я уверен, что ведение журнала работает, потому что я использовал его где-то еще.
- Чтобы исключить вероятность того, что это недавняя ошибка, я попробовал использовать Apache Beam 2.56 и 2.55.1, и у меня было та же проблема.
Подробнее здесь: https://stackoverflow.com/questions/784 ... d-dataflow