使用 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 的輕量情境 |