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]