[
https://issues.apache.org/jira/browse/FLINK-40198?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Dian Fu reassigned FLINK-40198:
-------------------------------
Assignee: niliushall
> Add Kafka source and sink APIs to DataFrame API
> -----------------------------------------------
>
> Key: FLINK-40198
> URL: https://issues.apache.org/jira/browse/FLINK-40198
> Project: Flink
> Issue Type: Sub-task
> Components: API / Python
> Reporter: Dian Fu
> Assignee: niliushall
> Priority: Major
> Fix For: 2.4.0
>
>
> 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: 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,
> ) -> DataFrame
> def DataFrame.write_kafka(
> self,
> bootstrap_servers: str,
> *,
> topic: str,
> format: str = "json",
> format_options: Optional[Dict[str, str]] = None,
> key_format: Optional[str] = None,
> key_fields: Optional[List[str]] = None,
> value_format: Optional[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,
> ) -> None
> read_kafka() and write_kafka() use Flink's Kafka SQL connector. The
> parameters map to Kafka connector options, including topic, topic-pattern,
> properties.*, scan.startup.*, scan.bounded.*, key/value format options,
> delivery guarantee, and topic partition discovery. Kafka header metadata can
> still be modeled through the existing connector metadata mechanism when
> needed.
> Example:
> events = pf.read_kafka(
> "localhost:9092",
> topic="events",
> format="json",
> schema={"user_id": DataType.int64(), "event": DataType.string()},
> startup_mode="earliest-offset",
> )
> events.write_kafka(
> "localhost:9092",
> topic="output-events",
> format="json",
> delivery_guarantee="at-least-once",
> )
--
This message was sent by Atlassian Jira
(v8.20.10#820010)