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]

Reply via email to