dianfu commented on code in PR #29245:
URL: https://github.com/apache/flink/pull/29245#discussion_r4069430565


##########
flink-python/pyflink/dataframe/dataframe.py:
##########
@@ -1369,6 +1369,205 @@ def columns(self) -> List[str]:
 
     # ======================== I/O ========================
 
+    @PublicEvolving()
+    def write_parquet(
+        self,
+        path: str,
+        *,
+        mode: str = "overwrite",

Review Comment:
   Do you think if it makes sense to set the `mode` to `append` for streaming 
jobs as this is the only supported mode for streaming jobs?  It's necessary to 
document this behavior clearly in the API doc.



##########
flink-python/pyflink/dataframe/dataframe.py:
##########
@@ -1369,6 +1369,205 @@ def columns(self) -> List[str]:
 
     # ======================== I/O ========================
 
+    @PublicEvolving()
+    def write_parquet(
+        self,
+        path: str,
+        *,
+        mode: str = "overwrite",
+        partition_by: Optional[Union[str, List[str]]] = None,
+        compression: Optional[str] = None,
+        utc_timezone: Optional[bool] = None,
+        sink_parallelism: Optional[int] = None,
+        rolling_policy_file_size: Optional[str] = None,
+        rolling_policy_rollover_interval: Optional[str] = None,
+        rolling_policy_inactivity_interval: Optional[str] = None,
+        rolling_policy_check_interval: Optional[str] = None,
+        partition_commit_trigger: Optional[str] = None,
+        partition_commit_delay: Optional[str] = None,
+        partition_commit_policy_kind: Optional[str] = None,

Review Comment:
   Could you also add the following options `sink.shuffle-by-partition.enable`, 
`auto-compaction`, `compaction.file-size` to the parameters of both 
write_parquet/write_json which I think may be frequently used?



##########
flink-python/pyflink/dataframe/io.py:
##########
@@ -16,15 +16,269 @@
 # limitations under the License.
 
################################################################################
 
-from typing import Dict, Optional, Tuple
+from typing import Dict, List, Optional, Tuple, Union
 
 from pyflink.dataframe.context import get_or_create_table_environment
-from pyflink.dataframe.dataframe import DataFrame
+from pyflink.dataframe.dataframe import DataFrame, _normalize_subset
 from pyflink.dataframe.datatype import DataType
 from pyflink.table import Schema, TableDescriptor
 from pyflink.util.api_stability_decorators import PublicEvolving
 
-__all__ = ["read_generic"]
+__all__ = ["read_generic", "read_json", "read_parquet"]
+
+
+def _build_filesystem_options(

Review Comment:
   The four option arguments represent two independent dimensions: connector 
vs. format options, and named convenience parameters vs. user-provided 
escape-hatch dictionaries. 
   
   The current names (`options`, `format_parameters`, `connector_options`, and 
`format_options`) do not make this distinction clear. 
   
   What about renaming them consistently, for example to 
`connector_parameters`, `format_parameters`, `extra_connector_options`, and 
`extra_format_options`, make them keyword-only.
   
   For example:
   ```
   def _build_filesystem_options(
       path: str,
       file_format: str,
       *,
       connector_parameters: Optional[Dict[str, Optional[str]]] = None,
       format_parameters: Optional[Dict[str, Optional[str]]] = None,
       extra_connector_options: Optional[Dict[str, str]] = None,
       extra_format_options: Optional[Dict[str, str]] = None,
   ) -> Dict[str, str]:
   ```



##########
flink-python/pyflink/dataframe/io.py:
##########
@@ -16,15 +16,269 @@
 # limitations under the License.
 
################################################################################
 
-from typing import Dict, Optional, Tuple
+from typing import Dict, List, Optional, Tuple, Union
 
 from pyflink.dataframe.context import get_or_create_table_environment
-from pyflink.dataframe.dataframe import DataFrame
+from pyflink.dataframe.dataframe import DataFrame, _normalize_subset
 from pyflink.dataframe.datatype import DataType
 from pyflink.table import Schema, TableDescriptor
 from pyflink.util.api_stability_decorators import PublicEvolving
 
-__all__ = ["read_generic"]
+__all__ = ["read_generic", "read_json", "read_parquet"]
+
+
+def _build_filesystem_options(
+    path: str,
+    file_format: str,
+    options: Dict[str, Optional[str]],
+    format_options: Optional[Dict[str, str]] = None,
+    *,
+    connector_options: Optional[Dict[str, str]] = None,
+    format_parameters: Optional[Dict[str, Optional[str]]] = None,
+) -> Dict[str, str]:
+    if not isinstance(path, str):
+        raise TypeError("path must be a string")
+    if not path:
+        raise ValueError("path must not be empty")
+
+    result = {"path": path, "format": file_format}
+    if connector_options is not None:
+        _validate_options(connector_options)
+        for key in connector_options:
+            if key in ("path", "format"):
+                raise ValueError(f"{key!r} must not be specified in 
connector_options")
+            if key.startswith(file_format + "."):
+                raise ValueError(f"format option {key!r} must be specified in 
format_options")
+        _merge_options(result, connector_options)
+    _merge_options(result, {key: value for key, value in options.items() if 
value is not None})
+
+    normalized_format_options: Dict[str, str] = {}
+    if format_options is not None:
+        _validate_options(format_options)
+        for key, value in format_options.items():
+            option = key if key.startswith(file_format + ".") else file_format 
+ "." + key
+            if option in normalized_format_options:
+                raise ValueError(f"duplicate format option: {option!r}")
+            normalized_format_options[option] = value
+    if format_parameters is not None:
+        _merge_options(normalized_format_options, {
+            file_format + "." + key: value
+            for key, value in format_parameters.items() if value is not None
+        })
+    _merge_options(result, normalized_format_options)
+    return result
+
+
+def _merge_options(target: Dict[str, str], options: Dict[str, str]) -> None:
+    _validate_options(options)
+    for key, value in options.items():
+        if key in target and target[key] != value:
+            raise ValueError(f"conflicting values for option {key!r}")
+        target[key] = value
+
+
+def _boolean_option(value: Optional[bool], name: str) -> Optional[str]:
+    if value is None:
+        return None
+    if not isinstance(value, bool):
+        raise TypeError(f"{name} must be a bool or None")
+    return str(value).lower()
+
+
+def _parallelism_option(value: Optional[int]) -> Optional[str]:
+    if value is None:
+        return None
+    if isinstance(value, bool) or not isinstance(value, int):
+        raise TypeError("sink_parallelism must be an int or None")
+    return str(value)
+
+
+def _build_filesystem_sink_options(
+    path: str,
+    file_format: str,
+    options: Dict[str, Optional[str]],
+    format_parameters: Dict[str, Optional[str]],
+    connector_options: Optional[Dict[str, str]] = None,
+    format_options: Optional[Dict[str, str]] = None,
+) -> Dict[str, str]:
+    result = _build_filesystem_options(
+        path, file_format, options, format_options,
+        connector_options=connector_options, 
format_parameters=format_parameters,
+    )
+    # Apply defaults after merging explicit settings from either API entry 
point.
+    defaults = {

Review Comment:
   This helper re-materializes only a subset of the upstream defaults, while 
other documented defaults are left to the connector factory. What about using 
one consistent model? Preferably keep `None` as “unspecified”, omit the option, 
and let the connector/format factory apply its default. The API docs can state 
the current effective default explicitly as following: `If None, the connector 
default is used (currently "SNAPPY")`.



##########
flink-python/pyflink/dataframe/dataframe.py:
##########
@@ -1369,6 +1369,205 @@ def columns(self) -> List[str]:
 
     # ======================== I/O ========================
 
+    @PublicEvolving()
+    def write_parquet(
+        self,
+        path: str,
+        *,
+        mode: str = "overwrite",
+        partition_by: Optional[Union[str, List[str]]] = None,
+        compression: Optional[str] = None,
+        utc_timezone: Optional[bool] = None,
+        sink_parallelism: Optional[int] = None,
+        rolling_policy_file_size: Optional[str] = None,
+        rolling_policy_rollover_interval: Optional[str] = None,
+        rolling_policy_inactivity_interval: Optional[str] = None,
+        rolling_policy_check_interval: Optional[str] = None,
+        partition_commit_trigger: Optional[str] = None,
+        partition_commit_delay: Optional[str] = None,
+        partition_commit_policy_kind: Optional[str] = None,
+        connector_options: Optional[Dict[str, str]] = None,
+        format_options: Optional[Dict[str, str]] = None,
+    ) -> None:
+        """
+        Write Parquet files using Flink's filesystem connector.
+
+        The filesystem connector and Parquet format must be available to 
Flink. The write
+        is submitted immediately and waits for completion for local or 
MiniCluster execution.
+        Sink columns are derived from this DataFrame's schema. Overwrite 
requires batch
+        execution; use ``mode="append"`` for streaming execution. Partitioned 
overwrite replaces
+        only partitions present in the input, retaining other partitions.
+
+        Rolling policies apply to streaming sinks. Parquet also rolls files on 
checkpoints;
+        continuous writes require checkpointing to finish files. Partition 
commit in streaming
+        requires ``partition_by`` and a commit policy. ``partition-time`` 
additionally requires
+        upstream watermarks and a partition time extractor, configured via 
``connector_options``.
+        For a TIMESTAMP_LTZ watermark, set 
``sink.partition-commit.watermark-time-zone`` in
+        ``connector_options`` to the session time zone; its default is UTC.
+
+        :param path: Output directory URI supported by Flink's filesystem 
implementations.
+        :param mode: ``"overwrite"`` (default) replaces existing data; 
``"append"`` adds files.
+        :param partition_by: Partition column name or non-empty list of names 
in directory order.
+            Values are stored in Hive-style partition paths rather than in the 
Parquet records.
+        :param compression: Parquet compression codec, defaulting to 
``"SNAPPY"``.
+        :param utc_timezone: Use UTC for Parquet timestamp conversion. 
Defaults to ``False``,
+            which uses the JVM default time zone, independently of the session 
time zone.
+        :param sink_parallelism: Sink parallelism; defaults to the upstream 
parallelism.
+        :param rolling_policy_file_size: Part file size threshold for rolling, 
default ``"128mb"``.
+            This is not a hard upper bound on file size.
+        :param rolling_policy_rollover_interval: Part file open-time 
threshold, default ``"30min"``.
+        :param rolling_policy_inactivity_interval: Part file inactivity 
threshold, default
+            ``"30min"``.
+        :param rolling_policy_check_interval: Interval for checking time-based 
rolling policies,
+            default ``"1min"``.
+        :param partition_commit_trigger: Partition commit trigger: 
``"process-time"`` (default) or
+            ``"partition-time"``.
+        :param partition_commit_delay: Delay before committing a partition, 
default ``"0s"``.
+        :param partition_commit_policy_kind: Optional comma-separated 
policies, such as
+            ``"success-file"`` or ``"custom"``. The ``metastore`` policy 
requires a Hive table.
+        :param connector_options: Additional filesystem options with string 
keys and values.
+            The ``connector``, ``path`` and ``format`` keys are reserved. 
Format options belong in
+            ``format_options``. Explicit parameter and dictionary values must 
agree when both
+            are set. Defaults are applied only after merging explicit settings.
+        :param format_options: Parquet options with string values, with or 
without the
+            ``parquet.`` prefix. Duplicate normalized keys are rejected. 
``None`` parameters
+            leave dictionary values unchanged; conflicting explicit values are 
rejected.
+        :raises TypeError: If an argument has an invalid type.
+        :raises ValueError: If the path or partition keys are invalid, the 
write mode is
+            unsupported, or options conflict.
+
+        Example::
+
+            >>> import pyflink.dataframe as pf
+            >>> _ = pf.config.set("execution.runtime-mode", "batch")
+            >>> events = pf.from_records([(1, "login")], schema=["id", 
"event"])
+            >>> events.write_parquet("file:///tmp/events", compression="GZIP")
+
+        .. versionadded:: 2.4.0
+        """
+        from pyflink.dataframe.io import (
+            _boolean_option,
+            _build_filesystem_sink_options,
+            _parallelism_option,
+        )
+
+        options = _build_filesystem_sink_options(
+            path,
+            "parquet",
+            {
+                "sink.parallelism": _parallelism_option(sink_parallelism),
+                "sink.rolling-policy.file-size": rolling_policy_file_size,
+                "sink.rolling-policy.rollover-interval": 
rolling_policy_rollover_interval,
+                "sink.rolling-policy.inactivity-interval": 
rolling_policy_inactivity_interval,
+                "sink.rolling-policy.check-interval": 
rolling_policy_check_interval,
+                "sink.partition-commit.trigger": partition_commit_trigger,
+                "sink.partition-commit.delay": partition_commit_delay,
+                "sink.partition-commit.policy.kind": 
partition_commit_policy_kind,
+            },
+            {
+                "compression": compression,
+                "utc-timezone": _boolean_option(utc_timezone, "utc_timezone"),
+            },
+            connector_options=connector_options,
+            format_options=format_options,
+        )
+        self._write("filesystem", options, mode, partition_by)
+
+    @PublicEvolving()
+    def write_json(
+        self,
+        path: str,
+        *,
+        mode: str = "overwrite",
+        partition_by: Optional[Union[str, List[str]]] = None,
+        timestamp_format: Optional[str] = None,
+        sink_parallelism: Optional[int] = None,
+        rolling_policy_file_size: Optional[str] = None,
+        rolling_policy_rollover_interval: Optional[str] = None,
+        rolling_policy_inactivity_interval: Optional[str] = None,
+        rolling_policy_check_interval: Optional[str] = None,
+        partition_commit_trigger: Optional[str] = None,
+        partition_commit_delay: Optional[str] = None,
+        partition_commit_policy_kind: Optional[str] = None,
+        connector_options: Optional[Dict[str, str]] = None,
+        format_options: Optional[Dict[str, str]] = None,

Review Comment:
   Could we also expose the following parameters `encode.ignore-null-fields` 
and `encode.decimal-as-plain-number` for write_json?



-- 
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