Dian Fu created FLINK-40198:
-------------------------------
Summary: 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
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)