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


##########
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:
   Added ignore_null_fields and decimal_as_plain_number to write_json, with 
tests checking the actual JSON output.



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