beetle0915 opened a new pull request, #29400: URL: https://github.com/apache/flink/pull/29400
## What is the purpose of the change Implements [FLINK-40429](https://issues.apache.org/jira/browse/FLINK-40429), following the multi-sink write design in [FLIP-591](https://cwiki.apache.org/confluence/spaces/FLINK/pages/430409184/FLIP-591+Introducing+Python+DataFrame+API+in+PyFlink). DataFrame writers currently submit each write immediately. This change lets users pass a shared StatementSet to stage multiple writes and submit them together, without manually converting DataFrames to Tables. Writes remain eager when no statement set is supplied. ```python sset = pf.create_statement_set() df.write_json("/tmp/events-json", statement_set=sset) df.write_catalog_table("events_sink", statement_set=sset) sset.execute() ``` ## Brief change log - Export `create_statement_set()` using the configured DataFrame environment and the existing Table API StatementSet. - Add the optional `statement_set` argument to `write_parquet`, `write_json`, `write_generic`, and `write_catalog_table`. Validate the set type and environment before staging inserts; preserve overwrite and partition options. - Document multi-sink execution and add unit and batch/streaming integration tests. ## Verifying this change - Added 10 tests covering factory behavior, staging without execution, argument/environment validation, overwrite and partition forwarding, explain behavior, error translation, and real JSON/catalog CSV multi-sink output in batch and streaming modes. - Re-ran 54 targeted tests after syncing with master: all passed. These include the new tests, context tests, existing generic/catalog/filesystem interface tests, and Table API StatementSet completeness. - An earlier broader DataFrame run on the same base passed 526 tests with 8 existing Hadoop/Parquet tests skipped. Seven native thread-mode tests were run separately and failed with JVM gateway disconnection; the same failures reproduced without this change. Their root cause remains unresolved. - Flake8, `git diff --check`, and a focused Sphinx build of the changed documentation with warnings treated as errors passed. - Full-repository `mvnw clean verify` and real Parquet multi-sink execution have not been run. Whole-site Sphinx validation was blocked by the existing `common/config.rst` reference to `pyflink.common.config_options`. This PR is a draft pending further validation, including community CI. ## Does this pull request potentially affect one of the following parts: - Dependencies: no - Public API: yes, adds the FLIP-591 factory and optional writer arguments - Serializers: no - Runtime per-record code paths: no - Deployment or recovery components: no - S3 file system connector implementation: no ## Documentation - Does this pull request introduce a new feature? Yes. - How is the feature documented? Python API docstrings and DataFrame environment/I/O reference documentation. --- ##### Was generative AI tooling used to co-author this PR? - [X] Yes Generated-by: OpenAI Codex (GPT-6) -- 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]
