Pandas DataFrames can be converted into a PyFlink Table.
Internally, PyFlink will serialize the Pandas DataFrame using Arrow columnar format on the client.
The serialized data will be processed and deserialized in Arrow source during execution.
The Arrow source can also be used in streaming jobs, and is integrated with checkpointing to
provide exactly-once guarantees.
The following example shows how to create a PyFlink Table from a Pandas DataFrame:
Convert PyFlink Table to Pandas DataFrame
PyFlink Tables can additionally be converted into a Pandas DataFrame.
The resulting rows will be serialized as multiple Arrow batches of Arrow columnar format on the client.
The maximum Arrow batch size is configured via the option python.fn-execution.arrow.batch.size.
The serialized data will then be converted to a Pandas DataFrame.
Because the contents of the table will be collected on the client, please ensure that the results of the table can fit in memory before calling this method.
You can limit the number of rows collected to client side via Table.limit
The following example shows how to convert a PyFlink Table to a Pandas DataFrame: