全部产品
Search
文档中心

数据湖构建:使用 PyPaimon 单机 SQL 查询 Paimon 表

更新时间:Jul 08, 2026

使用 PyPaimon 内置 SQL 引擎在单机 Python 环境直接查询 DLF Catalog 管理的 Paimon 表,无需部署 Flink/Spark 集群,适用于本地数据探查、模型调试、报表生成等 ad-hoc 场景。

前提条件

安装依赖

  1. 下载 PyPaimon 离线 wheel 包:pypaimon-1.5.dev20260601.tar.gz

  2. 安装离线包和其他依赖:

    pip install \
        pypaimon-1.5.dev20260601.tar.gz \
        datafusion \
        pypaimon-rust \
        pyarrow pandas requests
  3. 执行以下代码验证 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

固定为 rest

uri

DLF Paimon REST 访问端点,详见 服务接入点与公网访问。VPC 访问支持 HTTP 和 HTTPS 协议,公网访问必须使用 HTTPS 协议。

warehouse

DLF Catalog 名称。

token.provider

AccessKey 认证模式下固定为 dlf

dlf.region

DLF 所属地域 ID(如 cn-hangzhou),详见 服务接入点与公网访问

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 表(包含 idnamedeptsalary 四列)和 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 文档。

入口

表名引用

返回类型

适用场景

PaimonCatalog + datafusion.SessionContext(本文介绍)

三段式 paimon.<database>.<table>

DataFusion DataFrame,支持 .to_pylist().to_arrow_table().to_pandas() 等转换方法

需要与 DataFusion 生态衔接,或希望直接转换为 Arrow / Pandas / Polars 的场景

pypaimon.SQLContext

调用 set_current_catalogset_current_database 后可使用单段表名

list[pyarrow.RecordBatch],需要自行转换为 Pandas/Polars

仅需执行 SQL 并获取 RecordBatch 的轻量场景