使用 PyPaimon 内置 SQL 引擎在单机 Python 环境直接查询 DLF Catalog 管理的 Paimon 表,无需部署 Flink/Spark 集群,适用于本地数据探查、模型调试、报表生成等 ad-hoc 场景。
前提条件
安装依赖
下载 PyPaimon 离线 wheel 包:pypaimon-1.5.dev20260601.tar.gz
安装离线包和其他依赖:
pip install \ pypaimon-1.5.dev20260601.tar.gz \ datafusion \ pypaimon-rust \ pyarrow pandas requests执行以下代码验证 datafusion 与 PaimonCatalog 是否可用:
from datafusion import SessionContext from pypaimon_rust.datafusion import PaimonCatalog import datafusion print("datafusion:", datafusion.__version__) print("PaimonCatalog OK:", PaimonCatalog)
配置 Catalog 连接
在代码中定义 DLF Catalog 连接参数:
CATALOG_OPTIONS = {
"metastore": "rest",
"uri": "http://<DLF-ENDPOINT>",
"warehouse": "<YOUR-CATALOG>",
"token.provider": "dlf",
"dlf.region": "<REGION-ID>",
"dlf.access-key-id": "<ACCESS-KEY-ID>",
"dlf.access-key-secret": "<ACCESS-KEY-SECRET>",
"dlf.oss-endpoint": "<OSS-ENDPOINT>",
}各参数说明如下:
参数 | 说明 |
metastore | 固定为 |
uri | DLF Paimon REST 访问端点,详见 服务接入点与公网访问。VPC 访问支持 HTTP 和 HTTPS 协议,公网访问必须使用 HTTPS 协议。 |
warehouse | DLF Catalog 名称。 |
token.provider | AccessKey 认证模式下固定为 |
dlf.region | DLF 所属地域 ID(如 |
dlf.access-key-id / dlf.access-key-secret | 阿里云账号或 RAM 用户的 AccessKey。 |
dlf.oss-endpoint | (可选)OSS 公网访问 endpoint。公网访问必填,VPC 内访问可省略。 |
权限要求
使用 RAM 用户访问时,需要具备相应的 API 权限和数据权限,详情请参见快速配置权限。
注册 Paimon Catalog 到 DataFusion
PaimonCatalog 实现了 DataFusion 的 catalog provider 协议。注册到 SessionContext 之后,即可通过 paimon.<database>.<table> 三段式引用访问 Paimon 表。
from datafusion import SessionContext
from pypaimon_rust.datafusion import PaimonCatalog
ctx = SessionContext()
paimon = PaimonCatalog(CATALOG_OPTIONS)
ctx.register_catalog_provider("paimon", paimon)
print("catalogs:", ctx.catalog_names())
# -> {'paimon', 'datafusion'}
# 列出 paimon catalog 下所有 database(schema)
print("databases:", ctx.catalog("paimon").names())
# 列出 default database 下所有 table
print("tables in default:", ctx.catalog("paimon").schema("default").names())查询 Paimon 表
以下示例假设已存在 employees 表(包含 id、name、dept、salary 四列)和 departments 表。建表与写入代码见完整示例。
查询结果可以通过 .to_pylist()、.to_arrow_table()、.to_pandas() 等方法转换为下游需要的数据格式。
SELECT 查询
rows = ctx.sql("SELECT * FROM paimon.default.employees ORDER BY id").to_pylist()
print(rows[0])
# -> {'id': 1, 'name': 'Alice', 'dept': 'eng', 'salary': 9000}WHERE 过滤、列投影与 ORDER BY
eng = ctx.sql("""
SELECT id, name, salary
FROM paimon.default.employees
WHERE dept = 'eng'
ORDER BY salary DESC
""").to_pylist()
print(eng)
# -> [{'id': 4, 'name': 'Dan', 'salary': 9500},
# {'id': 1, 'name': 'Alice', 'salary': 9000},
# {'id': 2, 'name': 'Bob', 'salary': 8500}]GROUP BY 聚合
agg = ctx.sql("""
SELECT dept, COUNT(*) AS n, AVG(salary) AS avg_sal
FROM paimon.default.employees
GROUP BY dept
ORDER BY dept
""").to_pylist()
print(agg)
# -> [{'dept': 'eng', 'n': 3, 'avg_sal': 9000.0},
# {'dept': 'ops', 'n': 1, 'avg_sal': 6500.0},
# {'dept': 'sales', 'n': 2, 'avg_sal': 7100.0}]JOIN
自表 JOIN 示例:
result = ctx.sql("""
SELECT a.id AS id_a, b.id AS id_b, a.name
FROM paimon.default.employees a
JOIN paimon.default.employees b
ON a.id = b.id
WHERE a.dept = 'eng'
ORDER BY a.id
""").to_pylist()跨表 JOIN 示例:
result = ctx.sql("""
SELECT e.id, e.name, e.dept, d.dept_full
FROM paimon.default.employees e
JOIN paimon.default.departments d
ON e.dept = d.dept
WHERE e.dept = 'eng'
ORDER BY e.id
""").to_pylist()输出结果格式转换
查询结果可以转换为多种格式,衔接不同的下游处理框架。
转换为 Arrow Table
arrow_tbl = ctx.sql(
"SELECT id, salary FROM paimon.default.employees ORDER BY id"
).to_arrow_table()
print(arrow_tbl.num_rows, "rows")
print(arrow_tbl.schema)转换为 Pandas DataFrame
df = ctx.sql(
"SELECT * FROM paimon.default.employees ORDER BY id"
).to_pandas()
print(df.head())衔接 Daft
如果使用 Daft 框架,可以将 Arrow Table 转换为 Daft DataFrame 继续处理。关于 Daft 与 DLF Paimon 表的直接集成,请参见使用 Daft 操作 DLF Paimon 表。
import daft
daft_df = daft.from_arrow(arrow_tbl)
daft_df.show()完整示例
以下示例演示了端到端流程:使用 PyPaimon 建表并写入数据,然后注册到 DataFusion 并执行 SQL 查询。
from __future__ import annotations
import os
from datetime import datetime
import pyarrow as pa
from datafusion import SessionContext
from pypaimon_rust.datafusion import PaimonCatalog
from pypaimon import CatalogFactory, Schema
CATALOG_OPTIONS = {
"metastore": "rest",
"uri": "http://<DLF-ENDPOINT>",
"warehouse": "<YOUR-CATALOG>",
"token.provider": "dlf",
"dlf.region": "<REGION-ID>",
"dlf.access-key-id": os.environ["DLF_ACCESS_KEY_ID"],
"dlf.access-key-secret": os.environ["DLF_ACCESS_KEY_SECRET"],
"dlf.oss-endpoint": "<OSS-ENDPOINT>",
}
DATABASE = "default"
def main() -> None:
# 1) 使用 PyPaimon 建表并写入数据
pcat = CatalogFactory.create(CATALOG_OPTIONS)
pcat.create_database(DATABASE, True)
table_name = f"sql_demo_{datetime.now().strftime('%Y%m%d_%H%M%S')}"
full_name = f"{DATABASE}.{table_name}"
pa_schema = pa.schema([
pa.field("id", pa.int32(), nullable=False),
pa.field("name", pa.string()),
pa.field("dept", pa.string()),
pa.field("salary", pa.int64()),
])
schema = Schema.from_pyarrow_schema(
pa_schema,
options={"bucket": "-1", "file.format": "parquet"},
)
pcat.create_table(full_name, schema, True)
tbl = pcat.get_table(full_name)
wb = tbl.new_batch_write_builder()
w = wb.new_write()
w.write_arrow(pa.Table.from_pydict({
"id": [1, 2, 3, 4, 5, 6],
"name": ["Alice", "Bob", "Carol", "Dan", "Eve", "Frank"],
"dept": ["eng", "eng", "sales", "eng", "sales", "ops"],
"salary": [9000, 8500, 7000, 9500, 7200, 6500],
}, schema=pa_schema))
wb.new_commit().commit(w.prepare_commit())
w.close()
print(f"created + wrote 6 rows to {full_name}")
# 2) 注册 Catalog 到 DataFusion
ctx = SessionContext()
ctx.register_catalog_provider("paimon", PaimonCatalog(CATALOG_OPTIONS))
# 3) 执行 SQL
rows = ctx.sql(
f"SELECT * FROM paimon.{DATABASE}.{table_name} ORDER BY id"
).to_pylist()
print("all rows:", rows)
eng = ctx.sql(
f"SELECT id, name, salary FROM paimon.{DATABASE}.{table_name} "
f"WHERE dept = 'eng' ORDER BY salary DESC"
).to_pylist()
print("eng:", eng)
agg = ctx.sql(
f"SELECT dept, COUNT(*) AS n, AVG(salary) AS avg_sal "
f"FROM paimon.{DATABASE}.{table_name} GROUP BY dept ORDER BY dept"
).to_pylist()
print("by dept:", agg)
if __name__ == "__main__":
main()使用命令行入口
PyPaimon 自带 SQL CLI,通过 YAML 配置文件指定 Catalog 连接后,即可启动交互式 REPL 或执行单条 SQL。
将 Catalog 配置保存为 catalog.yaml:
metastore: rest
uri: http://<DLF-ENDPOINT>
warehouse: <YOUR-CATALOG>
token.provider: dlf
dlf.region: <REGION-ID>
dlf.access-key-id: <ACCESS-KEY-ID>
dlf.access-key-secret: <ACCESS-KEY-SECRET>
dlf.oss-endpoint: <OSS-ENDPOINT>通过 paimon 命令行调用:
# 启动交互式 REPL
paimon --config /path/to/catalog.yaml sql
# 执行单条 SQL
paimon --config /path/to/catalog.yaml sql "SELECT * FROM default.employees LIMIT 5"
# 指定输出格式(默认 table,可选 json)
paimon --config /path/to/catalog.yaml sql --format json "SELECT * FROM default.employees"更多选项(如输出格式、结果导出等)请执行 paimon sql --help 查看完整参数列表。
SQL 查询入口对比
PyPaimon 提供两种 SQL 查询入口,使用方式存在差异。本文介绍的是 PaimonCatalog + DataFusion 方式,SQLContext 方式请参见 PyPaimon API 文档。
入口 | 表名引用 | 返回类型 | 适用场景 |
| 三段式 | DataFusion DataFrame,支持 | 需要与 DataFusion 生态衔接,或希望直接转换为 Arrow / Pandas / Polars 的场景 |
| 调用 |
| 仅需执行 SQL 并获取 RecordBatch 的轻量场景 |