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]