romkabauer commented on issue #31830:
URL: https://github.com/apache/beam/issues/31830#issuecomment-2569802742
@ahmedabu98 Does it require additional adjustments to support streaming
mode? Or I'm just missing something?
While trying smth like
`iceberg_config = {
"table": "test.write_read",
"catalog_name": "default",
"catalog_properties": {
"type": "hadoop",
"warehouse": "file:///test_from_kafka_to_iceberg_write",
},
"triggering_frequency_seconds": 30
}
with beam.Pipeline(options=pipeline_options) as p:
(
p | "ReadExpensesFromKafka" >> ReadExpensesFromKafka(
bootstrap_servers=known_args.kafka_bootstrap,
topics=[known_args.kafka_input],
group_id="expenses-reader-group",
expansion_service=get_expansion_service()
)
| beam.WindowInto(
FixedWindows(1 * 30),
trigger=Repeatedly(AfterProcessingTime(1 * 30)),
accumulation_mode=AccumulationMode.DISCARDING
)
| 'Write to Iceberg' >> beam.managed.Write(
beam.managed.ICEBERG,
config=iceberg_config
)
)
`
I'm getting the following error - KeyError:
'after_synchronized_processing_time'
<details><summary>Traceback</summary>
<p>
Traceback (most recent call last):
File "/usr/local/lib/python3.10/runpy.py", line 196, in _run_module_as_main
return _run_code(code, main_globals, None,
File "/usr/local/lib/python3.10/runpy.py", line 86, in _run_code
exec(code, run_globals)
File "/app/beam_app/run.py", line 4, in <module>
run()
File "/app/beam_app/ingest_expenses_beam_app/pipeline_launcher.py", line
127, in run
p | "ReadExpensesFromKafka" >> ReadExpensesFromKafka(
File "/usr/local/lib/python3.10/site-packages/apache_beam/pvalue.py", line
138, in __or__
return self.pipeline.apply(ptransform, self)
File "/usr/local/lib/python3.10/site-packages/apache_beam/pipeline.py",
line 681, in apply
return self.apply(
File "/usr/local/lib/python3.10/site-packages/apache_beam/pipeline.py",
line 692, in apply
return self.apply(transform, pvalueish)
File "/usr/local/lib/python3.10/site-packages/apache_beam/pipeline.py",
line 754, in apply
pvalueish_result = self.runner.apply(transform, pvalueish, self._options)
File
"/usr/local/lib/python3.10/site-packages/apache_beam/runners/runner.py", line
191, in apply
return self.apply_PTransform(transform, input, options)
File
"/usr/local/lib/python3.10/site-packages/apache_beam/runners/runner.py", line
195, in apply_PTransform
return transform.expand(input)
File
"/usr/local/lib/python3.10/site-packages/apache_beam/transforms/managed.py",
line 162, in expand
return input | SchemaAwareExternalTransform(
File "/usr/local/lib/python3.10/site-packages/apache_beam/pvalue.py", line
138, in __or__
return self.pipeline.apply(ptransform, self)
File "/usr/local/lib/python3.10/site-packages/apache_beam/pipeline.py",
line 754, in apply
pvalueish_result = self.runner.apply(transform, pvalueish, self._options)
File
"/usr/local/lib/python3.10/site-packages/apache_beam/runners/runner.py", line
191, in apply
return self.apply_PTransform(transform, input, options)
File
"/usr/local/lib/python3.10/site-packages/apache_beam/runners/runner.py", line
195, in apply_PTransform
return transform.expand(input)
File
"/usr/local/lib/python3.10/site-packages/apache_beam/transforms/external.py",
line 429, in expand
return pcolls | self._payload_builder.identifier() >> ExternalTransform(
File "/usr/local/lib/python3.10/site-packages/apache_beam/pvalue.py", line
138, in __or__
return self.pipeline.apply(ptransform, self)
File "/usr/local/lib/python3.10/site-packages/apache_beam/pipeline.py",
line 681, in apply
return self.apply(
File "/usr/local/lib/python3.10/site-packages/apache_beam/pipeline.py",
line 692, in apply
return self.apply(transform, pvalueish)
File "/usr/local/lib/python3.10/site-packages/apache_beam/pipeline.py",
line 754, in apply
pvalueish_result = self.runner.apply(transform, pvalueish, self._options)
File
"/usr/local/lib/python3.10/site-packages/apache_beam/runners/runner.py", line
191, in apply
return self.apply_PTransform(transform, input, options)
File
"/usr/local/lib/python3.10/site-packages/apache_beam/runners/runner.py", line
195, in apply_PTransform
return transform.expand(input)
File
"/usr/local/lib/python3.10/site-packages/apache_beam/transforms/external.py",
line 773, in expand
self._outputs = {
File
"/usr/local/lib/python3.10/site-packages/apache_beam/transforms/external.py",
line 774, in <dictcomp>
tag: fix_output(result_context.pcollections.get_by_id(pcoll_id), tag)
File
"/usr/local/lib/python3.10/site-packages/apache_beam/runners/pipeline_context.py",
line 106, in get_by_id
self._id_to_obj[id] = self._obj_type.from_runner_api(
File "/usr/local/lib/python3.10/site-packages/apache_beam/pvalue.py", line
209, in from_runner_api
windowing=context.windowing_strategies.get_by_id(
File
"/usr/local/lib/python3.10/site-packages/apache_beam/runners/pipeline_context.py",
line 106, in get_by_id
self._id_to_obj[id] = self._obj_type.from_runner_api(
File
"/usr/local/lib/python3.10/site-packages/apache_beam/transforms/core.py", line
3709, in from_runner_api
triggerfn=TriggerFn.from_runner_api(proto.trigger, context),
File
"/usr/local/lib/python3.10/site-packages/apache_beam/transforms/trigger.py",
line 314, in from_runner_api
}[proto.WhichOneof('trigger')].from_runner_api(proto, context)
File
"/usr/local/lib/python3.10/site-packages/apache_beam/transforms/trigger.py",
line 734, in from_runner_api
TriggerFn.from_runner_api(proto.repeat.subtrigger, context))
File
"/usr/local/lib/python3.10/site-packages/apache_beam/transforms/trigger.py",
line 301, in from_runner_api
return {
KeyError: 'after_synchronized_processing_time'
</p>
</details>
As I understand after_synchronized_processing_time as trigger property was
deprecated some time ago?
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]