PyFlink Table 과 Pandas DataFrame 간의 변환
PyFlink Table 과 Pandas DataFrame 간의 변환
PyFlink Table API 는 PyFlink Table 과 Pandas DataFrame 사이의 변환을 지원합니다.
출처: 문서
본문
PyFlink Table API 는 PyFlink Table 과 Pandas DataFrame 사이의 변환을 지원합니다.
Pandas DataFrame 을 PyFlink Table 로 변환
Pandas DataFrame 은 PyFlink Table 로 변환될 수 있습니다. 내부적으로 PyFlink 는 클라이언트에서 Arrow columnar format 으로 Pandas DataFrame 을 직렬화합니다. 직렬화된 데이터는 실행 중에 Arrow source 에서 처리되고 역직렬화됩니다. Arrow source 는 스트리밍 작업에서도 사용할 수 있으며, 체크포인팅과 통합되어 exactly-once 보장을 제공합니다.
다음 예제는 Pandas DataFrame 에서 PyFlink Table 을 만드는 방법을 보여줍니다:
from pyflink.table import DataTypes
import pandas as pd
import numpy as np
# Create a Pandas DataFrame
pdf = pd.DataFrame(np.random.rand(1000, 2))
# Create a PyFlink Table from a Pandas DataFrame
table = t_env.from_pandas(pdf)
# Create a PyFlink Table from a Pandas DataFrame with the specified column names
table = t_env.from_pandas(pdf, ['f0', 'f1'])
# Create a PyFlink Table from a Pandas DataFrame with the specified column types
table = t_env.from_pandas(pdf, [DataTypes.DOUBLE(), DataTypes.DOUBLE()])
# Create a PyFlink Table from a Pandas DataFrame with the specified row type
table = t_env.from_pandas(pdf,
DataTypes.ROW([DataTypes.FIELD("f0", DataTypes.DOUBLE()),
DataTypes.FIELD("f1", DataTypes.DOUBLE())]))
PyFlink Table 을 Pandas DataFrame 으로 변환
PyFlink Table 은 추가로 Pandas DataFrame 으로 변환될 수 있습니다. 결과 행은 클라이언트에서 Arrow columnar format 의 여러 Arrow 배치로 직렬화됩니다. 최대 Arrow 배치 크기는 옵션 python.fn-execution.arrow.batch.size 로 구성됩니다. 직렬화된 데이터는 그 후 Pandas DataFrame 으로 변환됩니다. 테이블의 내용이 클라이언트에서 수집되므로, 이 메서드를 호출하기 전에 테이블의 결과가 메모리에 들어갈 수 있는지 확인하세요. Table.limit 를 사용해 클라이언트에 수집되는 행 수를 제한할 수 있습니다.
다음 예제는 PyFlink Table 을 Pandas DataFrame 으로 변환하는 방법을 보여줍니다:
from pyflink.table.expressions import col
import pandas as pd
import numpy as np
# Create a PyFlink Table
pdf = pd.DataFrame(np.random.rand(1000, 2))
table = t_env.from_pandas(pdf, ["a", "b"]).filter(col('a') > 0.5)
# Convert the PyFlink Table to a Pandas DataFrame
pdf = table.limit(100).to_pandas()