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]

Reply via email to