dianfu commented on code in PR #29439:
URL: https://github.com/apache/flink/pull/29439#discussion_r4236584517
##########
flink-python/pyflink/dataframe/io.py:
##########
@@ -276,6 +276,416 @@ def read_json(
)
+def _normalize_kafka_properties(
+ properties: Optional[Dict[str, str]], *, reserved: tuple
+) -> Optional[Dict[str, str]]:
+ if properties is None:
+ return None
+ _validate_options(properties)
+ normalized = {
+ key if key.startswith("properties.") else f"properties.{key}": value
+ for key, value in properties.items()
+ }
+ for key in reserved:
+ if key in normalized:
+ raise ValueError(
+ f"{key!r} must be configured through its dedicated DataFrame "
+ f"argument, not through properties")
+ return normalized
+
+
+def _merge_kafka_format_options(
+ options: Dict[str, str],
+ value_format: str,
+ *,
+ format_options: Optional[Dict[str, str]],
+ value_format_options: Optional[Dict[str, str]],
+) -> None:
+ normalized_format_options = _normalize_kafka_format_options(
+ format_options, value_format=value_format)
+ normalized_value_format_options = _normalize_kafka_format_options(
+ value_format_options, value_format=value_format)
+ _merge_options(options, normalized_format_options)
+ _merge_options(options, normalized_value_format_options)
+
+
+def _merge_kafka_key_format_options(
+ options: Dict[str, str],
+ key_format: str,
+ key_format_options: Optional[Dict[str, str]],
+) -> None:
+ if key_format_options is None:
+ return
+ _validate_options(key_format_options)
+ normalized = {
+ key if key.startswith("key.") else f"key.{key_format}.{key}": value
+ for key, value in key_format_options.items()
+ }
+ _merge_options(options, normalized)
+
+
+def _normalize_kafka_format_options(
+ options: Optional[Dict[str, str]],
+ *,
+ value_format: str,
+) -> Dict[str, str]:
+ if options is None:
+ return {}
+ _validate_options(options)
+ value_prefix = f"value.{value_format}."
+ normalized: Dict[str, str] = {}
+ for key, value in options.items():
+ if key.startswith(value_prefix):
+ normalized[key] = value
+ elif key.startswith(f"{value_format}."):
+ normalized[value_prefix + key[len(f"{value_format}."):]] = value
+ else:
+ normalized[value_prefix + key] = value
+ return normalized
+
+
+def _build_kafka_options(
+ bootstrap_servers: str,
+ *,
+ topic: Optional[Union[str, List[str]]],
+ topic_pattern: Optional[str],
+ group_id: Optional[str],
+ value_format: str,
+ format_options: Optional[Dict[str, str]],
+ value_format_options: Optional[Dict[str, str]],
+ key_format: Optional[str],
+ key_format_options: Optional[Dict[str, str]],
+ key_fields: Optional[List[str]],
+ key_fields_prefix: Optional[str],
+ value_fields_include: str,
+ startup_mode: str,
+ startup_specific_offsets: Optional[Union[str, Dict[int, int]]],
+ startup_timestamp_millis: Optional[int],
+ topic_partition_discovery_interval: Optional[str],
+ bounded_mode: str,
+ bounded_timestamp_millis: Optional[int],
+ bounded_specific_offsets: Optional[Union[str, Dict[int, int]]],
+ properties: Optional[Dict[str, str]],
+) -> Dict[str, str]:
+ if not isinstance(bootstrap_servers, str):
+ raise TypeError("bootstrap_servers must be a string")
+ if not bootstrap_servers:
+ raise ValueError("bootstrap_servers must not be empty")
+
+ if topic is None and topic_pattern is None:
+ raise ValueError("either 'topic' or 'topic_pattern' must be specified")
+ if topic is not None and topic_pattern is not None:
+ raise ValueError("'topic' and 'topic_pattern' are mutually exclusive")
+ if topic is not None:
+ if isinstance(topic, str):
+ resolved_topic = topic
+ elif isinstance(topic, list):
+ if not topic:
+ raise ValueError("topic must not be an empty list")
+ if not all(isinstance(t, str) for t in topic):
+ raise TypeError("topic list elements must be strings")
+ resolved_topic = ",".join(topic)
Review Comment:
I guess we should use semicolons
##########
flink-python/pyflink/dataframe/dataframe.py:
##########
@@ -2633,6 +2308,105 @@ def write_json(
statement_set=statement_set,
)
+ @PublicEvolving()
+ def write_kafka(
+ self,
+ bootstrap_servers: str,
+ *,
+ topic: str,
+ format: str = "json",
+ format_options: Optional[Dict[str, str]] = None,
+ key_format: Optional[str] = None,
+ key_format_options: Optional[Dict[str, str]] = None,
+ key_fields: Optional[List[str]] = None,
+ value_format: Optional[str] = None,
+ value_format_options: Optional[Dict[str, str]] = None,
+ value_fields_include: Literal["ALL", "EXCEPT_KEY"] = "ALL",
+ delivery_guarantee: Literal[
+ "none", "at-least-once", "exactly-once"
+ ] = "at-least-once",
Review Comment:
Missing option for `sink.transactional-id-prefix` which is needed when
delivery_guarantee is specified as `exactly-once`.
##########
flink-python/pyflink/dataframe/io.py:
##########
@@ -276,6 +276,416 @@ def read_json(
)
+def _normalize_kafka_properties(
+ properties: Optional[Dict[str, str]], *, reserved: tuple
+) -> Optional[Dict[str, str]]:
+ if properties is None:
+ return None
+ _validate_options(properties)
+ normalized = {
+ key if key.startswith("properties.") else f"properties.{key}": value
+ for key, value in properties.items()
+ }
+ for key in reserved:
+ if key in normalized:
+ raise ValueError(
+ f"{key!r} must be configured through its dedicated DataFrame "
+ f"argument, not through properties")
+ return normalized
+
+
+def _merge_kafka_format_options(
+ options: Dict[str, str],
+ value_format: str,
+ *,
+ format_options: Optional[Dict[str, str]],
+ value_format_options: Optional[Dict[str, str]],
+) -> None:
+ normalized_format_options = _normalize_kafka_format_options(
+ format_options, value_format=value_format)
+ normalized_value_format_options = _normalize_kafka_format_options(
+ value_format_options, value_format=value_format)
+ _merge_options(options, normalized_format_options)
+ _merge_options(options, normalized_value_format_options)
+
+
+def _merge_kafka_key_format_options(
+ options: Dict[str, str],
+ key_format: str,
+ key_format_options: Optional[Dict[str, str]],
+) -> None:
+ if key_format_options is None:
+ return
+ _validate_options(key_format_options)
+ normalized = {
+ key if key.startswith("key.") else f"key.{key_format}.{key}": value
+ for key, value in key_format_options.items()
+ }
+ _merge_options(options, normalized)
+
+
+def _normalize_kafka_format_options(
+ options: Optional[Dict[str, str]],
+ *,
+ value_format: str,
+) -> Dict[str, str]:
+ if options is None:
+ return {}
+ _validate_options(options)
+ value_prefix = f"value.{value_format}."
+ normalized: Dict[str, str] = {}
+ for key, value in options.items():
+ if key.startswith(value_prefix):
+ normalized[key] = value
+ elif key.startswith(f"{value_format}."):
+ normalized[value_prefix + key[len(f"{value_format}."):]] = value
+ else:
+ normalized[value_prefix + key] = value
+ return normalized
+
+
+def _build_kafka_options(
+ bootstrap_servers: str,
+ *,
+ topic: Optional[Union[str, List[str]]],
+ topic_pattern: Optional[str],
+ group_id: Optional[str],
+ value_format: str,
+ format_options: Optional[Dict[str, str]],
+ value_format_options: Optional[Dict[str, str]],
+ key_format: Optional[str],
+ key_format_options: Optional[Dict[str, str]],
+ key_fields: Optional[List[str]],
+ key_fields_prefix: Optional[str],
+ value_fields_include: str,
+ startup_mode: str,
+ startup_specific_offsets: Optional[Union[str, Dict[int, int]]],
+ startup_timestamp_millis: Optional[int],
+ topic_partition_discovery_interval: Optional[str],
+ bounded_mode: str,
+ bounded_timestamp_millis: Optional[int],
+ bounded_specific_offsets: Optional[Union[str, Dict[int, int]]],
+ properties: Optional[Dict[str, str]],
+) -> Dict[str, str]:
+ if not isinstance(bootstrap_servers, str):
+ raise TypeError("bootstrap_servers must be a string")
+ if not bootstrap_servers:
+ raise ValueError("bootstrap_servers must not be empty")
+
+ if topic is None and topic_pattern is None:
+ raise ValueError("either 'topic' or 'topic_pattern' must be specified")
+ if topic is not None and topic_pattern is not None:
+ raise ValueError("'topic' and 'topic_pattern' are mutually exclusive")
+ if topic is not None:
+ if isinstance(topic, str):
+ resolved_topic = topic
+ elif isinstance(topic, list):
+ if not topic:
+ raise ValueError("topic must not be an empty list")
+ if not all(isinstance(t, str) for t in topic):
+ raise TypeError("topic list elements must be strings")
+ resolved_topic = ",".join(topic)
+ else:
+ raise TypeError("topic must be a string or a list of strings")
+ else:
+ resolved_topic = None
+ if topic_pattern is not None and not isinstance(topic_pattern, str):
+ raise TypeError("topic_pattern must be a string")
+
+ if group_id is not None and not isinstance(group_id, str):
+ raise TypeError("group_id must be a string")
+
+ if not isinstance(value_format, str):
+ raise TypeError("format must be a string")
+ if not value_format:
+ raise ValueError("format must not be empty")
+
+ _validate_literal(value_fields_include, ("ALL", "EXCEPT_KEY"),
"value_fields_include")
+ _validate_literal(startup_mode, _STARTUP_MODES, "startup_mode")
+ _validate_literal(bounded_mode, _BOUNDED_MODES, "bounded_mode")
+
+ if startup_mode == "specific-offsets":
+ startup_specific_offsets = _convert_specific_offsets(
+ startup_specific_offsets, "startup_specific_offsets")
+ if startup_specific_offsets is None:
+ raise ValueError(
+ "startup_specific_offsets is required when "
+ "startup_mode='specific-offsets'")
+ else:
+ startup_specific_offsets = None
+ if startup_mode == "timestamp" and startup_timestamp_millis is None:
+ raise ValueError(
+ "startup_timestamp_millis is required when
startup_mode='timestamp'")
+ if (
+ startup_timestamp_millis is not None
+ and (
+ isinstance(startup_timestamp_millis, bool)
+ or not isinstance(startup_timestamp_millis, int)
+ )
+ ):
+ raise TypeError("startup_timestamp_millis must be an int")
+
+ if bounded_mode == "specific-offsets":
+ bounded_specific_offsets = _convert_specific_offsets(
+ bounded_specific_offsets, "bounded_specific_offsets")
+ if bounded_specific_offsets is None:
+ raise ValueError(
+ "bounded_specific_offsets is required when "
+ "bounded_mode='specific-offsets'")
+ else:
+ bounded_specific_offsets = None
+ if bounded_mode == "timestamp" and bounded_timestamp_millis is None:
+ raise ValueError(
+ "bounded_timestamp_millis is required when
bounded_mode='timestamp'")
+ if (
+ bounded_timestamp_millis is not None
+ and (
+ isinstance(bounded_timestamp_millis, bool)
+ or not isinstance(bounded_timestamp_millis, int)
+ )
+ ):
+ raise TypeError("bounded_timestamp_millis must be an int")
+
+ if (
+ topic_partition_discovery_interval is not None
+ and not isinstance(topic_partition_discovery_interval, str)
+ ):
+ raise TypeError("topic_partition_discovery_interval must be a string")
+
+ if key_format is not None and not isinstance(key_format, str):
+ raise TypeError("key_format must be a string")
+ if key_fields is not None:
+ if not isinstance(key_fields, list):
+ raise TypeError("key_fields must be a list of strings")
+ if not key_fields:
+ raise ValueError("key_fields must not be empty")
+ if not all(isinstance(f, str) for f in key_fields):
+ raise TypeError("key_fields elements must be strings")
+ if key_fields_prefix is not None and not isinstance(key_fields_prefix,
str):
+ raise TypeError("key_fields_prefix must be a string")
+
+ normalized_properties = _normalize_kafka_properties(
+ properties, reserved=("properties.bootstrap.servers",
"properties.group.id"))
+
+ options: Dict[str, str] = {
+ "properties.bootstrap.servers": bootstrap_servers,
+ "value.format": value_format,
+ "value.fields-include": value_fields_include,
+ }
+ if group_id is not None:
+ options["properties.group.id"] = group_id
+ if resolved_topic is not None:
+ options["topic"] = resolved_topic
+ if topic_pattern is not None:
+ options["topic-pattern"] = topic_pattern
+ if topic_partition_discovery_interval is not None:
+ options["scan.topic-partition-discovery.interval"] = (
+ topic_partition_discovery_interval)
+ options["scan.startup.mode"] = startup_mode
+ if startup_specific_offsets is not None:
+ options["scan.startup.specific-offsets"] = startup_specific_offsets
+ if startup_timestamp_millis is not None:
+ options["scan.startup.timestamp-millis"] =
str(startup_timestamp_millis)
+ if bounded_mode != "unbounded":
+ options["scan.bounded.mode"] = bounded_mode
+ if bounded_timestamp_millis is not None:
+ options["scan.bounded.timestamp-millis"] =
str(bounded_timestamp_millis)
+ if bounded_specific_offsets is not None:
+ options["scan.bounded.specific-offsets"] = bounded_specific_offsets
+
+ if key_format is not None:
+ options["key.format"] = key_format
+ _merge_kafka_key_format_options(options, key_format, key_format_options)
+ if key_fields is not None:
+ options["key.fields"] = ",".join(key_fields)
+ if key_fields_prefix is not None:
+ options["key.fields-prefix"] = key_fields_prefix
+
+ _merge_kafka_format_options(
+ options,
+ value_format,
+ format_options=format_options,
+ value_format_options=value_format_options,
+ )
+
+ if normalized_properties is not None:
+ _merge_options(options, normalized_properties)
+
+ return options
+
+
+_STARTUP_MODES = (
+ "earliest-offset", "latest-offset", "group-offsets",
+ "timestamp", "specific-offsets")
+_BOUNDED_MODES = (
+ "unbounded", "group-offsets", "latest-offset",
+ "timestamp", "specific-offsets")
+_DELIVERY_GUARANTEES = ("none", "at-least-once", "exactly-once")
+
+
+def _validate_literal(value: str, choices: tuple, name: str) -> None:
+ if not isinstance(value, str):
+ raise TypeError(f"{name} must be a string")
+ if value not in choices:
+ raise ValueError(
+ f"{name} must be one of {choices}, got {value!r}")
+
+
+def _convert_specific_offsets(
+ value: Optional[Union[str, Dict[int, int]]], name: str
+) -> Optional[str]:
+ if value is None:
+ return None
+ if isinstance(value, str):
+ return value
+ if isinstance(value, dict):
+ parts = []
+ for partition, offset in value.items():
+ if not isinstance(partition, int) or isinstance(partition, bool):
+ raise TypeError(f"{name} keys must be ints (partition
numbers)")
+ if not isinstance(offset, int) or isinstance(offset, bool):
+ raise TypeError(f"{name} values must be ints (offsets)")
+ parts.append(f"{partition}:{offset}")
+ return ",".join(parts)
Review Comment:
I guess we should use semicolons.
##########
flink-python/pyflink/dataframe/io.py:
##########
@@ -276,6 +276,416 @@ def read_json(
)
+def _normalize_kafka_properties(
+ properties: Optional[Dict[str, str]], *, reserved: tuple
+) -> Optional[Dict[str, str]]:
+ if properties is None:
+ return None
+ _validate_options(properties)
+ normalized = {
+ key if key.startswith("properties.") else f"properties.{key}": value
+ for key, value in properties.items()
+ }
+ for key in reserved:
+ if key in normalized:
+ raise ValueError(
+ f"{key!r} must be configured through its dedicated DataFrame "
+ f"argument, not through properties")
+ return normalized
+
+
+def _merge_kafka_format_options(
+ options: Dict[str, str],
+ value_format: str,
+ *,
+ format_options: Optional[Dict[str, str]],
+ value_format_options: Optional[Dict[str, str]],
+) -> None:
+ normalized_format_options = _normalize_kafka_format_options(
+ format_options, value_format=value_format)
+ normalized_value_format_options = _normalize_kafka_format_options(
+ value_format_options, value_format=value_format)
+ _merge_options(options, normalized_format_options)
+ _merge_options(options, normalized_value_format_options)
+
+
+def _merge_kafka_key_format_options(
+ options: Dict[str, str],
+ key_format: str,
+ key_format_options: Optional[Dict[str, str]],
+) -> None:
+ if key_format_options is None:
+ return
+ _validate_options(key_format_options)
+ normalized = {
+ key if key.startswith("key.") else f"key.{key_format}.{key}": value
+ for key, value in key_format_options.items()
+ }
+ _merge_options(options, normalized)
+
+
+def _normalize_kafka_format_options(
+ options: Optional[Dict[str, str]],
+ *,
+ value_format: str,
+) -> Dict[str, str]:
+ if options is None:
+ return {}
+ _validate_options(options)
+ value_prefix = f"value.{value_format}."
+ normalized: Dict[str, str] = {}
+ for key, value in options.items():
+ if key.startswith(value_prefix):
+ normalized[key] = value
+ elif key.startswith(f"{value_format}."):
+ normalized[value_prefix + key[len(f"{value_format}."):]] = value
+ else:
+ normalized[value_prefix + key] = value
+ return normalized
+
+
+def _build_kafka_options(
+ bootstrap_servers: str,
+ *,
+ topic: Optional[Union[str, List[str]]],
+ topic_pattern: Optional[str],
+ group_id: Optional[str],
+ value_format: str,
+ format_options: Optional[Dict[str, str]],
+ value_format_options: Optional[Dict[str, str]],
+ key_format: Optional[str],
+ key_format_options: Optional[Dict[str, str]],
+ key_fields: Optional[List[str]],
+ key_fields_prefix: Optional[str],
+ value_fields_include: str,
+ startup_mode: str,
+ startup_specific_offsets: Optional[Union[str, Dict[int, int]]],
+ startup_timestamp_millis: Optional[int],
+ topic_partition_discovery_interval: Optional[str],
+ bounded_mode: str,
+ bounded_timestamp_millis: Optional[int],
+ bounded_specific_offsets: Optional[Union[str, Dict[int, int]]],
+ properties: Optional[Dict[str, str]],
+) -> Dict[str, str]:
+ if not isinstance(bootstrap_servers, str):
+ raise TypeError("bootstrap_servers must be a string")
+ if not bootstrap_servers:
+ raise ValueError("bootstrap_servers must not be empty")
+
+ if topic is None and topic_pattern is None:
+ raise ValueError("either 'topic' or 'topic_pattern' must be specified")
+ if topic is not None and topic_pattern is not None:
+ raise ValueError("'topic' and 'topic_pattern' are mutually exclusive")
+ if topic is not None:
+ if isinstance(topic, str):
+ resolved_topic = topic
+ elif isinstance(topic, list):
+ if not topic:
+ raise ValueError("topic must not be an empty list")
+ if not all(isinstance(t, str) for t in topic):
+ raise TypeError("topic list elements must be strings")
+ resolved_topic = ",".join(topic)
+ else:
+ raise TypeError("topic must be a string or a list of strings")
+ else:
+ resolved_topic = None
+ if topic_pattern is not None and not isinstance(topic_pattern, str):
+ raise TypeError("topic_pattern must be a string")
+
+ if group_id is not None and not isinstance(group_id, str):
+ raise TypeError("group_id must be a string")
+
+ if not isinstance(value_format, str):
+ raise TypeError("format must be a string")
+ if not value_format:
+ raise ValueError("format must not be empty")
+
+ _validate_literal(value_fields_include, ("ALL", "EXCEPT_KEY"),
"value_fields_include")
+ _validate_literal(startup_mode, _STARTUP_MODES, "startup_mode")
+ _validate_literal(bounded_mode, _BOUNDED_MODES, "bounded_mode")
+
+ if startup_mode == "specific-offsets":
+ startup_specific_offsets = _convert_specific_offsets(
+ startup_specific_offsets, "startup_specific_offsets")
+ if startup_specific_offsets is None:
+ raise ValueError(
+ "startup_specific_offsets is required when "
+ "startup_mode='specific-offsets'")
+ else:
+ startup_specific_offsets = None
+ if startup_mode == "timestamp" and startup_timestamp_millis is None:
+ raise ValueError(
+ "startup_timestamp_millis is required when
startup_mode='timestamp'")
+ if (
+ startup_timestamp_millis is not None
+ and (
+ isinstance(startup_timestamp_millis, bool)
+ or not isinstance(startup_timestamp_millis, int)
+ )
+ ):
+ raise TypeError("startup_timestamp_millis must be an int")
+
+ if bounded_mode == "specific-offsets":
+ bounded_specific_offsets = _convert_specific_offsets(
+ bounded_specific_offsets, "bounded_specific_offsets")
+ if bounded_specific_offsets is None:
+ raise ValueError(
+ "bounded_specific_offsets is required when "
+ "bounded_mode='specific-offsets'")
+ else:
+ bounded_specific_offsets = None
+ if bounded_mode == "timestamp" and bounded_timestamp_millis is None:
+ raise ValueError(
+ "bounded_timestamp_millis is required when
bounded_mode='timestamp'")
+ if (
+ bounded_timestamp_millis is not None
+ and (
+ isinstance(bounded_timestamp_millis, bool)
+ or not isinstance(bounded_timestamp_millis, int)
+ )
+ ):
+ raise TypeError("bounded_timestamp_millis must be an int")
+
+ if (
+ topic_partition_discovery_interval is not None
+ and not isinstance(topic_partition_discovery_interval, str)
+ ):
+ raise TypeError("topic_partition_discovery_interval must be a string")
+
+ if key_format is not None and not isinstance(key_format, str):
+ raise TypeError("key_format must be a string")
+ if key_fields is not None:
+ if not isinstance(key_fields, list):
+ raise TypeError("key_fields must be a list of strings")
+ if not key_fields:
+ raise ValueError("key_fields must not be empty")
+ if not all(isinstance(f, str) for f in key_fields):
+ raise TypeError("key_fields elements must be strings")
+ if key_fields_prefix is not None and not isinstance(key_fields_prefix,
str):
+ raise TypeError("key_fields_prefix must be a string")
+
+ normalized_properties = _normalize_kafka_properties(
+ properties, reserved=("properties.bootstrap.servers",
"properties.group.id"))
+
+ options: Dict[str, str] = {
+ "properties.bootstrap.servers": bootstrap_servers,
+ "value.format": value_format,
+ "value.fields-include": value_fields_include,
+ }
+ if group_id is not None:
+ options["properties.group.id"] = group_id
+ if resolved_topic is not None:
+ options["topic"] = resolved_topic
+ if topic_pattern is not None:
+ options["topic-pattern"] = topic_pattern
+ if topic_partition_discovery_interval is not None:
+ options["scan.topic-partition-discovery.interval"] = (
+ topic_partition_discovery_interval)
+ options["scan.startup.mode"] = startup_mode
+ if startup_specific_offsets is not None:
+ options["scan.startup.specific-offsets"] = startup_specific_offsets
+ if startup_timestamp_millis is not None:
+ options["scan.startup.timestamp-millis"] =
str(startup_timestamp_millis)
+ if bounded_mode != "unbounded":
+ options["scan.bounded.mode"] = bounded_mode
+ if bounded_timestamp_millis is not None:
+ options["scan.bounded.timestamp-millis"] =
str(bounded_timestamp_millis)
+ if bounded_specific_offsets is not None:
+ options["scan.bounded.specific-offsets"] = bounded_specific_offsets
+
+ if key_format is not None:
+ options["key.format"] = key_format
+ _merge_kafka_key_format_options(options, key_format, key_format_options)
+ if key_fields is not None:
+ options["key.fields"] = ",".join(key_fields)
+ if key_fields_prefix is not None:
+ options["key.fields-prefix"] = key_fields_prefix
+
+ _merge_kafka_format_options(
+ options,
+ value_format,
+ format_options=format_options,
+ value_format_options=value_format_options,
+ )
+
+ if normalized_properties is not None:
+ _merge_options(options, normalized_properties)
+
+ return options
+
+
+_STARTUP_MODES = (
+ "earliest-offset", "latest-offset", "group-offsets",
+ "timestamp", "specific-offsets")
+_BOUNDED_MODES = (
+ "unbounded", "group-offsets", "latest-offset",
+ "timestamp", "specific-offsets")
+_DELIVERY_GUARANTEES = ("none", "at-least-once", "exactly-once")
+
+
+def _validate_literal(value: str, choices: tuple, name: str) -> None:
+ if not isinstance(value, str):
+ raise TypeError(f"{name} must be a string")
+ if value not in choices:
+ raise ValueError(
+ f"{name} must be one of {choices}, got {value!r}")
+
+
+def _convert_specific_offsets(
+ value: Optional[Union[str, Dict[int, int]]], name: str
+) -> Optional[str]:
+ if value is None:
+ return None
+ if isinstance(value, str):
+ return value
+ if isinstance(value, dict):
+ parts = []
+ for partition, offset in value.items():
+ if not isinstance(partition, int) or isinstance(partition, bool):
+ raise TypeError(f"{name} keys must be ints (partition
numbers)")
+ if not isinstance(offset, int) or isinstance(offset, bool):
+ raise TypeError(f"{name} values must be ints (offsets)")
+ parts.append(f"{partition}:{offset}")
+ return ",".join(parts)
+ raise TypeError(f"{name} must be a str or a Dict[int, int]")
+
+
+@PublicEvolving()
+def read_kafka(
+ bootstrap_servers: str,
+ *,
+ topic: Optional[Union[str, List[str]]] = None,
+ topic_pattern: Optional[str] = None,
+ group_id: Optional[str] = None,
+ schema: Dict[str, DataType],
Review Comment:
Move `schema` to right after `*` to keep it consistent with the other APIs.
##########
flink-python/pyflink/dataframe/io.py:
##########
@@ -276,6 +276,416 @@ def read_json(
)
+def _normalize_kafka_properties(
+ properties: Optional[Dict[str, str]], *, reserved: tuple
+) -> Optional[Dict[str, str]]:
+ if properties is None:
+ return None
+ _validate_options(properties)
+ normalized = {
+ key if key.startswith("properties.") else f"properties.{key}": value
+ for key, value in properties.items()
+ }
+ for key in reserved:
+ if key in normalized:
+ raise ValueError(
+ f"{key!r} must be configured through its dedicated DataFrame "
+ f"argument, not through properties")
+ return normalized
+
+
+def _merge_kafka_format_options(
+ options: Dict[str, str],
+ value_format: str,
+ *,
+ format_options: Optional[Dict[str, str]],
+ value_format_options: Optional[Dict[str, str]],
+) -> None:
+ normalized_format_options = _normalize_kafka_format_options(
+ format_options, value_format=value_format)
+ normalized_value_format_options = _normalize_kafka_format_options(
+ value_format_options, value_format=value_format)
+ _merge_options(options, normalized_format_options)
+ _merge_options(options, normalized_value_format_options)
+
+
+def _merge_kafka_key_format_options(
+ options: Dict[str, str],
+ key_format: str,
+ key_format_options: Optional[Dict[str, str]],
+) -> None:
+ if key_format_options is None:
+ return
+ _validate_options(key_format_options)
+ normalized = {
+ key if key.startswith("key.") else f"key.{key_format}.{key}": value
+ for key, value in key_format_options.items()
+ }
+ _merge_options(options, normalized)
+
+
+def _normalize_kafka_format_options(
+ options: Optional[Dict[str, str]],
+ *,
+ value_format: str,
+) -> Dict[str, str]:
+ if options is None:
+ return {}
+ _validate_options(options)
+ value_prefix = f"value.{value_format}."
+ normalized: Dict[str, str] = {}
+ for key, value in options.items():
+ if key.startswith(value_prefix):
+ normalized[key] = value
+ elif key.startswith(f"{value_format}."):
+ normalized[value_prefix + key[len(f"{value_format}."):]] = value
+ else:
+ normalized[value_prefix + key] = value
+ return normalized
+
+
+def _build_kafka_options(
+ bootstrap_servers: str,
+ *,
+ topic: Optional[Union[str, List[str]]],
+ topic_pattern: Optional[str],
+ group_id: Optional[str],
+ value_format: str,
+ format_options: Optional[Dict[str, str]],
+ value_format_options: Optional[Dict[str, str]],
+ key_format: Optional[str],
+ key_format_options: Optional[Dict[str, str]],
+ key_fields: Optional[List[str]],
+ key_fields_prefix: Optional[str],
+ value_fields_include: str,
+ startup_mode: str,
+ startup_specific_offsets: Optional[Union[str, Dict[int, int]]],
+ startup_timestamp_millis: Optional[int],
+ topic_partition_discovery_interval: Optional[str],
+ bounded_mode: str,
+ bounded_timestamp_millis: Optional[int],
+ bounded_specific_offsets: Optional[Union[str, Dict[int, int]]],
+ properties: Optional[Dict[str, str]],
+) -> Dict[str, str]:
+ if not isinstance(bootstrap_servers, str):
+ raise TypeError("bootstrap_servers must be a string")
+ if not bootstrap_servers:
+ raise ValueError("bootstrap_servers must not be empty")
+
+ if topic is None and topic_pattern is None:
+ raise ValueError("either 'topic' or 'topic_pattern' must be specified")
+ if topic is not None and topic_pattern is not None:
+ raise ValueError("'topic' and 'topic_pattern' are mutually exclusive")
+ if topic is not None:
+ if isinstance(topic, str):
+ resolved_topic = topic
+ elif isinstance(topic, list):
+ if not topic:
+ raise ValueError("topic must not be an empty list")
+ if not all(isinstance(t, str) for t in topic):
+ raise TypeError("topic list elements must be strings")
+ resolved_topic = ",".join(topic)
+ else:
+ raise TypeError("topic must be a string or a list of strings")
+ else:
+ resolved_topic = None
+ if topic_pattern is not None and not isinstance(topic_pattern, str):
+ raise TypeError("topic_pattern must be a string")
+
+ if group_id is not None and not isinstance(group_id, str):
+ raise TypeError("group_id must be a string")
+
+ if not isinstance(value_format, str):
+ raise TypeError("format must be a string")
+ if not value_format:
+ raise ValueError("format must not be empty")
+
+ _validate_literal(value_fields_include, ("ALL", "EXCEPT_KEY"),
"value_fields_include")
+ _validate_literal(startup_mode, _STARTUP_MODES, "startup_mode")
+ _validate_literal(bounded_mode, _BOUNDED_MODES, "bounded_mode")
+
+ if startup_mode == "specific-offsets":
+ startup_specific_offsets = _convert_specific_offsets(
+ startup_specific_offsets, "startup_specific_offsets")
+ if startup_specific_offsets is None:
+ raise ValueError(
+ "startup_specific_offsets is required when "
+ "startup_mode='specific-offsets'")
+ else:
+ startup_specific_offsets = None
+ if startup_mode == "timestamp" and startup_timestamp_millis is None:
+ raise ValueError(
+ "startup_timestamp_millis is required when
startup_mode='timestamp'")
+ if (
+ startup_timestamp_millis is not None
+ and (
+ isinstance(startup_timestamp_millis, bool)
+ or not isinstance(startup_timestamp_millis, int)
+ )
+ ):
+ raise TypeError("startup_timestamp_millis must be an int")
+
+ if bounded_mode == "specific-offsets":
+ bounded_specific_offsets = _convert_specific_offsets(
+ bounded_specific_offsets, "bounded_specific_offsets")
+ if bounded_specific_offsets is None:
+ raise ValueError(
+ "bounded_specific_offsets is required when "
+ "bounded_mode='specific-offsets'")
+ else:
+ bounded_specific_offsets = None
+ if bounded_mode == "timestamp" and bounded_timestamp_millis is None:
+ raise ValueError(
+ "bounded_timestamp_millis is required when
bounded_mode='timestamp'")
+ if (
+ bounded_timestamp_millis is not None
+ and (
+ isinstance(bounded_timestamp_millis, bool)
+ or not isinstance(bounded_timestamp_millis, int)
+ )
+ ):
+ raise TypeError("bounded_timestamp_millis must be an int")
+
+ if (
+ topic_partition_discovery_interval is not None
+ and not isinstance(topic_partition_discovery_interval, str)
+ ):
+ raise TypeError("topic_partition_discovery_interval must be a string")
+
+ if key_format is not None and not isinstance(key_format, str):
+ raise TypeError("key_format must be a string")
+ if key_fields is not None:
+ if not isinstance(key_fields, list):
+ raise TypeError("key_fields must be a list of strings")
+ if not key_fields:
+ raise ValueError("key_fields must not be empty")
+ if not all(isinstance(f, str) for f in key_fields):
+ raise TypeError("key_fields elements must be strings")
+ if key_fields_prefix is not None and not isinstance(key_fields_prefix,
str):
+ raise TypeError("key_fields_prefix must be a string")
+
+ normalized_properties = _normalize_kafka_properties(
+ properties, reserved=("properties.bootstrap.servers",
"properties.group.id"))
+
+ options: Dict[str, str] = {
+ "properties.bootstrap.servers": bootstrap_servers,
+ "value.format": value_format,
+ "value.fields-include": value_fields_include,
+ }
+ if group_id is not None:
+ options["properties.group.id"] = group_id
+ if resolved_topic is not None:
+ options["topic"] = resolved_topic
+ if topic_pattern is not None:
+ options["topic-pattern"] = topic_pattern
+ if topic_partition_discovery_interval is not None:
+ options["scan.topic-partition-discovery.interval"] = (
+ topic_partition_discovery_interval)
+ options["scan.startup.mode"] = startup_mode
+ if startup_specific_offsets is not None:
+ options["scan.startup.specific-offsets"] = startup_specific_offsets
+ if startup_timestamp_millis is not None:
+ options["scan.startup.timestamp-millis"] =
str(startup_timestamp_millis)
+ if bounded_mode != "unbounded":
+ options["scan.bounded.mode"] = bounded_mode
+ if bounded_timestamp_millis is not None:
+ options["scan.bounded.timestamp-millis"] =
str(bounded_timestamp_millis)
+ if bounded_specific_offsets is not None:
+ options["scan.bounded.specific-offsets"] = bounded_specific_offsets
+
+ if key_format is not None:
+ options["key.format"] = key_format
+ _merge_kafka_key_format_options(options, key_format, key_format_options)
+ if key_fields is not None:
+ options["key.fields"] = ",".join(key_fields)
+ if key_fields_prefix is not None:
+ options["key.fields-prefix"] = key_fields_prefix
+
+ _merge_kafka_format_options(
+ options,
+ value_format,
+ format_options=format_options,
+ value_format_options=value_format_options,
+ )
+
+ if normalized_properties is not None:
+ _merge_options(options, normalized_properties)
+
+ return options
+
+
+_STARTUP_MODES = (
+ "earliest-offset", "latest-offset", "group-offsets",
+ "timestamp", "specific-offsets")
+_BOUNDED_MODES = (
+ "unbounded", "group-offsets", "latest-offset",
+ "timestamp", "specific-offsets")
+_DELIVERY_GUARANTEES = ("none", "at-least-once", "exactly-once")
+
+
+def _validate_literal(value: str, choices: tuple, name: str) -> None:
+ if not isinstance(value, str):
+ raise TypeError(f"{name} must be a string")
+ if value not in choices:
+ raise ValueError(
+ f"{name} must be one of {choices}, got {value!r}")
+
+
+def _convert_specific_offsets(
+ value: Optional[Union[str, Dict[int, int]]], name: str
+) -> Optional[str]:
+ if value is None:
+ return None
+ if isinstance(value, str):
+ return value
+ if isinstance(value, dict):
+ parts = []
+ for partition, offset in value.items():
+ if not isinstance(partition, int) or isinstance(partition, bool):
+ raise TypeError(f"{name} keys must be ints (partition
numbers)")
+ if not isinstance(offset, int) or isinstance(offset, bool):
+ raise TypeError(f"{name} values must be ints (offsets)")
+ parts.append(f"{partition}:{offset}")
+ return ",".join(parts)
+ raise TypeError(f"{name} must be a str or a Dict[int, int]")
+
+
+@PublicEvolving()
+def read_kafka(
+ bootstrap_servers: str,
+ *,
+ topic: Optional[Union[str, List[str]]] = None,
+ topic_pattern: Optional[str] = None,
+ group_id: Optional[str] = None,
+ schema: Dict[str, DataType],
+ format: str = "json",
+ format_options: Optional[Dict[str, str]] = None,
+ key_format: Optional[str] = None,
+ key_format_options: Optional[Dict[str, str]] = None,
+ key_fields: Optional[List[str]] = None,
+ key_fields_prefix: Optional[str] = None,
+ value_format: Optional[str] = None,
+ value_format_options: Optional[Dict[str, str]] = None,
+ value_fields_include: Literal["ALL", "EXCEPT_KEY"] = "ALL",
+ startup_mode: Literal[
+ "earliest-offset", "latest-offset", "group-offsets",
+ "timestamp", "specific-offsets"
+ ] = "group-offsets",
+ startup_specific_offsets: Optional[Union[str, Dict[int, int]]] = None,
+ startup_timestamp_millis: Optional[int] = None,
+ topic_partition_discovery_interval: Optional[str] = "5 min",
+ bounded_mode: Literal[
+ "unbounded", "group-offsets", "latest-offset",
+ "timestamp", "specific-offsets"
+ ] = "unbounded",
+ bounded_timestamp_millis: Optional[int] = None,
+ bounded_specific_offsets: Optional[Union[str, Dict[int, int]]] = None,
Review Comment:
Do you think it makes sense to introduce parameter `options` to let users
specify connector options via key/value pairs. Besides, it's also useful in
case it introduces new options in the Kafka connector in the future. We are not
possible to make sure they are sync always.
Same for the write_kafka API.
##########
flink-python/pyflink/dataframe/io.py:
##########
@@ -276,6 +276,416 @@ def read_json(
)
+def _normalize_kafka_properties(
+ properties: Optional[Dict[str, str]], *, reserved: tuple
+) -> Optional[Dict[str, str]]:
+ if properties is None:
+ return None
+ _validate_options(properties)
+ normalized = {
+ key if key.startswith("properties.") else f"properties.{key}": value
+ for key, value in properties.items()
+ }
+ for key in reserved:
+ if key in normalized:
+ raise ValueError(
+ f"{key!r} must be configured through its dedicated DataFrame "
+ f"argument, not through properties")
+ return normalized
+
+
+def _merge_kafka_format_options(
+ options: Dict[str, str],
+ value_format: str,
+ *,
+ format_options: Optional[Dict[str, str]],
+ value_format_options: Optional[Dict[str, str]],
+) -> None:
+ normalized_format_options = _normalize_kafka_format_options(
+ format_options, value_format=value_format)
+ normalized_value_format_options = _normalize_kafka_format_options(
+ value_format_options, value_format=value_format)
+ _merge_options(options, normalized_format_options)
+ _merge_options(options, normalized_value_format_options)
+
+
+def _merge_kafka_key_format_options(
+ options: Dict[str, str],
+ key_format: str,
+ key_format_options: Optional[Dict[str, str]],
+) -> None:
+ if key_format_options is None:
+ return
+ _validate_options(key_format_options)
+ normalized = {
+ key if key.startswith("key.") else f"key.{key_format}.{key}": value
+ for key, value in key_format_options.items()
+ }
+ _merge_options(options, normalized)
+
+
+def _normalize_kafka_format_options(
+ options: Optional[Dict[str, str]],
+ *,
+ value_format: str,
+) -> Dict[str, str]:
+ if options is None:
+ return {}
+ _validate_options(options)
+ value_prefix = f"value.{value_format}."
+ normalized: Dict[str, str] = {}
+ for key, value in options.items():
+ if key.startswith(value_prefix):
+ normalized[key] = value
+ elif key.startswith(f"{value_format}."):
+ normalized[value_prefix + key[len(f"{value_format}."):]] = value
+ else:
+ normalized[value_prefix + key] = value
+ return normalized
+
+
+def _build_kafka_options(
+ bootstrap_servers: str,
+ *,
+ topic: Optional[Union[str, List[str]]],
+ topic_pattern: Optional[str],
+ group_id: Optional[str],
+ value_format: str,
+ format_options: Optional[Dict[str, str]],
+ value_format_options: Optional[Dict[str, str]],
+ key_format: Optional[str],
+ key_format_options: Optional[Dict[str, str]],
+ key_fields: Optional[List[str]],
+ key_fields_prefix: Optional[str],
+ value_fields_include: str,
+ startup_mode: str,
+ startup_specific_offsets: Optional[Union[str, Dict[int, int]]],
+ startup_timestamp_millis: Optional[int],
+ topic_partition_discovery_interval: Optional[str],
+ bounded_mode: str,
+ bounded_timestamp_millis: Optional[int],
+ bounded_specific_offsets: Optional[Union[str, Dict[int, int]]],
+ properties: Optional[Dict[str, str]],
+) -> Dict[str, str]:
+ if not isinstance(bootstrap_servers, str):
+ raise TypeError("bootstrap_servers must be a string")
+ if not bootstrap_servers:
+ raise ValueError("bootstrap_servers must not be empty")
+
+ if topic is None and topic_pattern is None:
+ raise ValueError("either 'topic' or 'topic_pattern' must be specified")
+ if topic is not None and topic_pattern is not None:
+ raise ValueError("'topic' and 'topic_pattern' are mutually exclusive")
+ if topic is not None:
+ if isinstance(topic, str):
+ resolved_topic = topic
+ elif isinstance(topic, list):
+ if not topic:
+ raise ValueError("topic must not be an empty list")
+ if not all(isinstance(t, str) for t in topic):
+ raise TypeError("topic list elements must be strings")
+ resolved_topic = ",".join(topic)
+ else:
+ raise TypeError("topic must be a string or a list of strings")
+ else:
+ resolved_topic = None
+ if topic_pattern is not None and not isinstance(topic_pattern, str):
+ raise TypeError("topic_pattern must be a string")
+
+ if group_id is not None and not isinstance(group_id, str):
+ raise TypeError("group_id must be a string")
+
+ if not isinstance(value_format, str):
+ raise TypeError("format must be a string")
+ if not value_format:
+ raise ValueError("format must not be empty")
+
+ _validate_literal(value_fields_include, ("ALL", "EXCEPT_KEY"),
"value_fields_include")
+ _validate_literal(startup_mode, _STARTUP_MODES, "startup_mode")
+ _validate_literal(bounded_mode, _BOUNDED_MODES, "bounded_mode")
+
+ if startup_mode == "specific-offsets":
+ startup_specific_offsets = _convert_specific_offsets(
+ startup_specific_offsets, "startup_specific_offsets")
+ if startup_specific_offsets is None:
+ raise ValueError(
+ "startup_specific_offsets is required when "
+ "startup_mode='specific-offsets'")
+ else:
+ startup_specific_offsets = None
+ if startup_mode == "timestamp" and startup_timestamp_millis is None:
+ raise ValueError(
+ "startup_timestamp_millis is required when
startup_mode='timestamp'")
+ if (
+ startup_timestamp_millis is not None
+ and (
+ isinstance(startup_timestamp_millis, bool)
+ or not isinstance(startup_timestamp_millis, int)
+ )
+ ):
+ raise TypeError("startup_timestamp_millis must be an int")
+
+ if bounded_mode == "specific-offsets":
+ bounded_specific_offsets = _convert_specific_offsets(
+ bounded_specific_offsets, "bounded_specific_offsets")
+ if bounded_specific_offsets is None:
+ raise ValueError(
+ "bounded_specific_offsets is required when "
+ "bounded_mode='specific-offsets'")
+ else:
+ bounded_specific_offsets = None
+ if bounded_mode == "timestamp" and bounded_timestamp_millis is None:
+ raise ValueError(
+ "bounded_timestamp_millis is required when
bounded_mode='timestamp'")
+ if (
+ bounded_timestamp_millis is not None
+ and (
+ isinstance(bounded_timestamp_millis, bool)
+ or not isinstance(bounded_timestamp_millis, int)
+ )
+ ):
+ raise TypeError("bounded_timestamp_millis must be an int")
+
+ if (
+ topic_partition_discovery_interval is not None
+ and not isinstance(topic_partition_discovery_interval, str)
+ ):
+ raise TypeError("topic_partition_discovery_interval must be a string")
+
+ if key_format is not None and not isinstance(key_format, str):
+ raise TypeError("key_format must be a string")
+ if key_fields is not None:
+ if not isinstance(key_fields, list):
+ raise TypeError("key_fields must be a list of strings")
+ if not key_fields:
+ raise ValueError("key_fields must not be empty")
+ if not all(isinstance(f, str) for f in key_fields):
+ raise TypeError("key_fields elements must be strings")
+ if key_fields_prefix is not None and not isinstance(key_fields_prefix,
str):
+ raise TypeError("key_fields_prefix must be a string")
+
+ normalized_properties = _normalize_kafka_properties(
+ properties, reserved=("properties.bootstrap.servers",
"properties.group.id"))
+
+ options: Dict[str, str] = {
+ "properties.bootstrap.servers": bootstrap_servers,
+ "value.format": value_format,
+ "value.fields-include": value_fields_include,
+ }
+ if group_id is not None:
+ options["properties.group.id"] = group_id
+ if resolved_topic is not None:
+ options["topic"] = resolved_topic
+ if topic_pattern is not None:
+ options["topic-pattern"] = topic_pattern
+ if topic_partition_discovery_interval is not None:
+ options["scan.topic-partition-discovery.interval"] = (
+ topic_partition_discovery_interval)
+ options["scan.startup.mode"] = startup_mode
+ if startup_specific_offsets is not None:
+ options["scan.startup.specific-offsets"] = startup_specific_offsets
+ if startup_timestamp_millis is not None:
+ options["scan.startup.timestamp-millis"] =
str(startup_timestamp_millis)
+ if bounded_mode != "unbounded":
+ options["scan.bounded.mode"] = bounded_mode
+ if bounded_timestamp_millis is not None:
+ options["scan.bounded.timestamp-millis"] =
str(bounded_timestamp_millis)
+ if bounded_specific_offsets is not None:
+ options["scan.bounded.specific-offsets"] = bounded_specific_offsets
+
+ if key_format is not None:
+ options["key.format"] = key_format
+ _merge_kafka_key_format_options(options, key_format, key_format_options)
+ if key_fields is not None:
+ options["key.fields"] = ",".join(key_fields)
+ if key_fields_prefix is not None:
+ options["key.fields-prefix"] = key_fields_prefix
+
+ _merge_kafka_format_options(
+ options,
+ value_format,
+ format_options=format_options,
+ value_format_options=value_format_options,
+ )
+
+ if normalized_properties is not None:
+ _merge_options(options, normalized_properties)
+
+ return options
+
+
+_STARTUP_MODES = (
+ "earliest-offset", "latest-offset", "group-offsets",
+ "timestamp", "specific-offsets")
+_BOUNDED_MODES = (
+ "unbounded", "group-offsets", "latest-offset",
+ "timestamp", "specific-offsets")
+_DELIVERY_GUARANTEES = ("none", "at-least-once", "exactly-once")
+
+
+def _validate_literal(value: str, choices: tuple, name: str) -> None:
+ if not isinstance(value, str):
+ raise TypeError(f"{name} must be a string")
+ if value not in choices:
+ raise ValueError(
+ f"{name} must be one of {choices}, got {value!r}")
+
+
+def _convert_specific_offsets(
+ value: Optional[Union[str, Dict[int, int]]], name: str
+) -> Optional[str]:
+ if value is None:
+ return None
+ if isinstance(value, str):
+ return value
+ if isinstance(value, dict):
+ parts = []
+ for partition, offset in value.items():
+ if not isinstance(partition, int) or isinstance(partition, bool):
+ raise TypeError(f"{name} keys must be ints (partition
numbers)")
+ if not isinstance(offset, int) or isinstance(offset, bool):
+ raise TypeError(f"{name} values must be ints (offsets)")
+ parts.append(f"{partition}:{offset}")
+ return ",".join(parts)
+ raise TypeError(f"{name} must be a str or a Dict[int, int]")
+
+
+@PublicEvolving()
+def read_kafka(
+ bootstrap_servers: str,
+ *,
+ topic: Optional[Union[str, List[str]]] = None,
+ topic_pattern: Optional[str] = None,
+ group_id: Optional[str] = None,
+ schema: Dict[str, DataType],
+ format: str = "json",
+ format_options: Optional[Dict[str, str]] = None,
+ key_format: Optional[str] = None,
+ key_format_options: Optional[Dict[str, str]] = None,
+ key_fields: Optional[List[str]] = None,
+ key_fields_prefix: Optional[str] = None,
+ value_format: Optional[str] = None,
+ value_format_options: Optional[Dict[str, str]] = None,
+ value_fields_include: Literal["ALL", "EXCEPT_KEY"] = "ALL",
+ startup_mode: Literal[
+ "earliest-offset", "latest-offset", "group-offsets",
+ "timestamp", "specific-offsets"
+ ] = "group-offsets",
+ startup_specific_offsets: Optional[Union[str, Dict[int, int]]] = None,
+ startup_timestamp_millis: Optional[int] = None,
+ topic_partition_discovery_interval: Optional[str] = "5 min",
+ bounded_mode: Literal[
+ "unbounded", "group-offsets", "latest-offset",
+ "timestamp", "specific-offsets"
+ ] = "unbounded",
+ bounded_timestamp_millis: Optional[int] = None,
+ bounded_specific_offsets: Optional[Union[str, Dict[int, int]]] = None,
+ properties: Optional[Dict[str, str]] = None,
+ computed_columns: Optional[Dict[str, str]] = None,
+ watermark: Optional[Tuple[str, str]] = None,
+) -> DataFrame:
+ """
+ Read data from Kafka using Flink's Kafka SQL connector.
+
+ The Kafka connector must be available to Flink. Exactly one of ``topic`` or
+ ``topic_pattern`` must be specified. Physical columns in ``schema`` are
followed by
+ computed columns in dictionary insertion order. A watermark can reference
a physical
+ or computed timestamp column.
+
+ ``format`` and ``value_format`` are aliases; ``value_format`` takes
precedence when
+ both are set.
+ ``format_options`` and ``value_format_options`` are both normalized to
+ ``value.<format>.<key>``; conflicting values are rejected.
+
+ :param bootstrap_servers: Comma-separated Kafka bootstrap server addresses.
+ :param topic: Kafka topic name or list of names. Mutually exclusive with
+ ``topic_pattern``.
+ :param topic_pattern: Regular expression matching Kafka topic names.
Mutually
+ exclusive with ``topic``.
+ :param group_id: Kafka consumer group id. Maps to ``properties.group.id``.
+ :param schema: Non-empty mapping of physical column names to DataFrame
data types.
+ :param format: Value format, for example ``"json"``, ``"csv"``, ``"avro"``.
+ :param format_options: Value format options with string values, with or
without
+ the ``<format>.`` prefix.
+ :param key_format: Key format, for example ``"json"`` or ``"avro"``.
+ :param key_format_options: Key format options with string values, with or
without
+ the ``key.<key_format>.`` prefix.
+ :param key_fields: List of column names that make up the key.
+ :param key_fields_prefix: Prefix for key fields to avoid name clashes with
value
+ fields.
+ :param value_format: Alternative to ``format`` for specifying the value
format.
+ :param value_format_options: Alternative to ``format_options`` for
specifying
+ value format options.
+ :param value_fields_include: Whether value fields include key fields:
+ ``"ALL"`` or ``"EXCEPT_KEY"``.
+ :param startup_mode: Startup mode: ``"earliest-offset"``,
``"latest-offset"``,
+ ``"group-offsets"``, ``"timestamp"``, or ``"specific-offsets"``.
+ :param startup_specific_offsets: Specific offsets for ``startup_mode`` of
+ ``"specific-offsets"``, as a string (``"partition:offset,..."``) or a
dict
+ (``{partition: offset}``).
+ :param startup_timestamp_millis: Startup timestamp in milliseconds for
+ ``startup_mode`` of ``"timestamp"``.
+ :param topic_partition_discovery_interval: Interval for topic partition
Review Comment:
If None means disabled, Should set scan.topic-partition-discovery.interval
to 0 when None is specified.
##########
flink-python/pyflink/dataframe/dataframe.py:
##########
@@ -2633,6 +2308,105 @@ def write_json(
statement_set=statement_set,
)
+ @PublicEvolving()
+ def write_kafka(
+ self,
+ bootstrap_servers: str,
+ *,
+ topic: str,
+ format: str = "json",
+ format_options: Optional[Dict[str, str]] = None,
+ key_format: Optional[str] = None,
+ key_format_options: Optional[Dict[str, str]] = None,
+ key_fields: Optional[List[str]] = None,
+ value_format: Optional[str] = None,
+ value_format_options: Optional[Dict[str, str]] = None,
+ value_fields_include: Literal["ALL", "EXCEPT_KEY"] = "ALL",
+ delivery_guarantee: Literal[
+ "none", "at-least-once", "exactly-once"
+ ] = "at-least-once",
+ properties: Optional[Dict[str, str]] = None,
+ parallelism: Optional[int] = None,
+ statement_set: Optional[StatementSet] = None,
Review Comment:
Missing option `sink.partitioner`
##########
flink-python/pyflink/dataframe/io.py:
##########
@@ -276,6 +276,416 @@ def read_json(
)
+def _normalize_kafka_properties(
+ properties: Optional[Dict[str, str]], *, reserved: tuple
+) -> Optional[Dict[str, str]]:
+ if properties is None:
+ return None
+ _validate_options(properties)
+ normalized = {
+ key if key.startswith("properties.") else f"properties.{key}": value
+ for key, value in properties.items()
+ }
+ for key in reserved:
+ if key in normalized:
+ raise ValueError(
+ f"{key!r} must be configured through its dedicated DataFrame "
+ f"argument, not through properties")
+ return normalized
+
+
+def _merge_kafka_format_options(
+ options: Dict[str, str],
+ value_format: str,
+ *,
+ format_options: Optional[Dict[str, str]],
+ value_format_options: Optional[Dict[str, str]],
+) -> None:
+ normalized_format_options = _normalize_kafka_format_options(
+ format_options, value_format=value_format)
+ normalized_value_format_options = _normalize_kafka_format_options(
+ value_format_options, value_format=value_format)
+ _merge_options(options, normalized_format_options)
+ _merge_options(options, normalized_value_format_options)
+
+
+def _merge_kafka_key_format_options(
+ options: Dict[str, str],
+ key_format: str,
+ key_format_options: Optional[Dict[str, str]],
+) -> None:
+ if key_format_options is None:
+ return
+ _validate_options(key_format_options)
+ normalized = {
+ key if key.startswith("key.") else f"key.{key_format}.{key}": value
+ for key, value in key_format_options.items()
+ }
+ _merge_options(options, normalized)
+
+
+def _normalize_kafka_format_options(
+ options: Optional[Dict[str, str]],
+ *,
+ value_format: str,
+) -> Dict[str, str]:
+ if options is None:
+ return {}
+ _validate_options(options)
+ value_prefix = f"value.{value_format}."
+ normalized: Dict[str, str] = {}
+ for key, value in options.items():
+ if key.startswith(value_prefix):
+ normalized[key] = value
+ elif key.startswith(f"{value_format}."):
+ normalized[value_prefix + key[len(f"{value_format}."):]] = value
+ else:
+ normalized[value_prefix + key] = value
+ return normalized
+
+
+def _build_kafka_options(
+ bootstrap_servers: str,
+ *,
+ topic: Optional[Union[str, List[str]]],
+ topic_pattern: Optional[str],
+ group_id: Optional[str],
+ value_format: str,
+ format_options: Optional[Dict[str, str]],
+ value_format_options: Optional[Dict[str, str]],
+ key_format: Optional[str],
+ key_format_options: Optional[Dict[str, str]],
+ key_fields: Optional[List[str]],
+ key_fields_prefix: Optional[str],
+ value_fields_include: str,
+ startup_mode: str,
+ startup_specific_offsets: Optional[Union[str, Dict[int, int]]],
+ startup_timestamp_millis: Optional[int],
+ topic_partition_discovery_interval: Optional[str],
+ bounded_mode: str,
+ bounded_timestamp_millis: Optional[int],
+ bounded_specific_offsets: Optional[Union[str, Dict[int, int]]],
+ properties: Optional[Dict[str, str]],
+) -> Dict[str, str]:
+ if not isinstance(bootstrap_servers, str):
+ raise TypeError("bootstrap_servers must be a string")
+ if not bootstrap_servers:
+ raise ValueError("bootstrap_servers must not be empty")
+
+ if topic is None and topic_pattern is None:
+ raise ValueError("either 'topic' or 'topic_pattern' must be specified")
+ if topic is not None and topic_pattern is not None:
+ raise ValueError("'topic' and 'topic_pattern' are mutually exclusive")
+ if topic is not None:
+ if isinstance(topic, str):
+ resolved_topic = topic
+ elif isinstance(topic, list):
+ if not topic:
+ raise ValueError("topic must not be an empty list")
+ if not all(isinstance(t, str) for t in topic):
+ raise TypeError("topic list elements must be strings")
+ resolved_topic = ",".join(topic)
+ else:
+ raise TypeError("topic must be a string or a list of strings")
+ else:
+ resolved_topic = None
+ if topic_pattern is not None and not isinstance(topic_pattern, str):
+ raise TypeError("topic_pattern must be a string")
+
+ if group_id is not None and not isinstance(group_id, str):
+ raise TypeError("group_id must be a string")
+
+ if not isinstance(value_format, str):
+ raise TypeError("format must be a string")
+ if not value_format:
+ raise ValueError("format must not be empty")
+
+ _validate_literal(value_fields_include, ("ALL", "EXCEPT_KEY"),
"value_fields_include")
+ _validate_literal(startup_mode, _STARTUP_MODES, "startup_mode")
+ _validate_literal(bounded_mode, _BOUNDED_MODES, "bounded_mode")
+
+ if startup_mode == "specific-offsets":
+ startup_specific_offsets = _convert_specific_offsets(
+ startup_specific_offsets, "startup_specific_offsets")
+ if startup_specific_offsets is None:
+ raise ValueError(
+ "startup_specific_offsets is required when "
+ "startup_mode='specific-offsets'")
+ else:
+ startup_specific_offsets = None
+ if startup_mode == "timestamp" and startup_timestamp_millis is None:
+ raise ValueError(
+ "startup_timestamp_millis is required when
startup_mode='timestamp'")
+ if (
+ startup_timestamp_millis is not None
+ and (
+ isinstance(startup_timestamp_millis, bool)
+ or not isinstance(startup_timestamp_millis, int)
+ )
+ ):
+ raise TypeError("startup_timestamp_millis must be an int")
+
+ if bounded_mode == "specific-offsets":
+ bounded_specific_offsets = _convert_specific_offsets(
+ bounded_specific_offsets, "bounded_specific_offsets")
+ if bounded_specific_offsets is None:
+ raise ValueError(
+ "bounded_specific_offsets is required when "
+ "bounded_mode='specific-offsets'")
+ else:
+ bounded_specific_offsets = None
+ if bounded_mode == "timestamp" and bounded_timestamp_millis is None:
+ raise ValueError(
+ "bounded_timestamp_millis is required when
bounded_mode='timestamp'")
+ if (
+ bounded_timestamp_millis is not None
+ and (
+ isinstance(bounded_timestamp_millis, bool)
+ or not isinstance(bounded_timestamp_millis, int)
+ )
+ ):
+ raise TypeError("bounded_timestamp_millis must be an int")
+
+ if (
+ topic_partition_discovery_interval is not None
+ and not isinstance(topic_partition_discovery_interval, str)
+ ):
+ raise TypeError("topic_partition_discovery_interval must be a string")
+
+ if key_format is not None and not isinstance(key_format, str):
+ raise TypeError("key_format must be a string")
+ if key_fields is not None:
+ if not isinstance(key_fields, list):
+ raise TypeError("key_fields must be a list of strings")
+ if not key_fields:
+ raise ValueError("key_fields must not be empty")
+ if not all(isinstance(f, str) for f in key_fields):
+ raise TypeError("key_fields elements must be strings")
+ if key_fields_prefix is not None and not isinstance(key_fields_prefix,
str):
+ raise TypeError("key_fields_prefix must be a string")
+
+ normalized_properties = _normalize_kafka_properties(
+ properties, reserved=("properties.bootstrap.servers",
"properties.group.id"))
+
+ options: Dict[str, str] = {
+ "properties.bootstrap.servers": bootstrap_servers,
+ "value.format": value_format,
+ "value.fields-include": value_fields_include,
+ }
+ if group_id is not None:
+ options["properties.group.id"] = group_id
+ if resolved_topic is not None:
+ options["topic"] = resolved_topic
+ if topic_pattern is not None:
+ options["topic-pattern"] = topic_pattern
+ if topic_partition_discovery_interval is not None:
+ options["scan.topic-partition-discovery.interval"] = (
+ topic_partition_discovery_interval)
+ options["scan.startup.mode"] = startup_mode
+ if startup_specific_offsets is not None:
+ options["scan.startup.specific-offsets"] = startup_specific_offsets
+ if startup_timestamp_millis is not None:
+ options["scan.startup.timestamp-millis"] =
str(startup_timestamp_millis)
+ if bounded_mode != "unbounded":
+ options["scan.bounded.mode"] = bounded_mode
+ if bounded_timestamp_millis is not None:
+ options["scan.bounded.timestamp-millis"] =
str(bounded_timestamp_millis)
+ if bounded_specific_offsets is not None:
+ options["scan.bounded.specific-offsets"] = bounded_specific_offsets
+
+ if key_format is not None:
+ options["key.format"] = key_format
+ _merge_kafka_key_format_options(options, key_format, key_format_options)
+ if key_fields is not None:
+ options["key.fields"] = ",".join(key_fields)
Review Comment:
I guess we should use semicolons. This issue also exists in the write_kafka
method.
##########
flink-python/pyflink/dataframe/dataframe.py:
##########
@@ -33,6 +32,7 @@
Type,
TypeVar,
Union,
+ Literal,
Review Comment:
already imported
--
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]