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)

Reply via email to