bingqin2 commented on issue #72982: URL: https://github.com/apache/airflow/issues/72982#issuecomment-5915147881
Hi @garvit-arora, go ahead, it's yours. A few pointers from when I wrote this up: - The other `DataflowConfiguration` fields become pipeline options in `BeamDataflowMixin.__get_dataflow_pipeline_options` in `providers/apache/beam/src/airflow/providers/apache/beam/operators/beam.py`. Keys there are camelCase (`serviceAccount`, `impersonateServiceAccount`); the Python and Go operators convert them to snake_case (`_init_pipeline_options(format_pipeline_options=True)`), the Java operator passes them as they are. So `maxNumWorkers`, set only when `max_num_workers` is not `None`, should reach all three SDKs. - `_init_pipeline_options` applies the operator's own `pipeline_options` after the Dataflow ones, so an explicit `max_num_workers` / `maxNumWorkers` pipeline option would win. Worth stating in the docstring and covering in a test. - Unit tests live in `providers/apache/beam/tests/unit/apache/beam/operators/test_beam.py`. The native Python system tests (`providers/google/tests/system/google/cloud/dataflow/example_dataflow_native_python.py`) already pass `max_num_workers: 1` through `dataflow_config`, so today they silently ignore it; after the fix they exercise it. -- 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]
