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]

Reply via email to