[ 
https://issues.apache.org/jira/browse/FLINK-40198?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=18118728#comment-18118728
 ] 

niliushall commented on FLINK-40198:
------------------------------------

Hello [~dianfu] , Could you assign this task to me?

> 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
>            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)

Reply via email to