Metadata-Version: 2.4
Name: hologres-dataframe
Version: 0.1.0
Summary: Lazy, pandas-style DataFrame API for Hologres
Requires-Python: >=3.9
Description-Content-Type: text/markdown
Requires-Dist: hologres-client>=0.1.2.dev0
Requires-Dist: psycopg[binary,pool]>=3.1
Requires-Dist: pandas>=1.3
Requires-Dist: pyarrow>=12.0
Requires-Dist: alibabacloud-fc20230330<5,>=4.7.9
Provides-Extra: fc
Provides-Extra: dev
Requires-Dist: pytest>=7.0; extra == "dev"
Requires-Dist: pytest-cov>=4.0; extra == "dev"
Requires-Dist: ruff>=0.4; extra == "dev"
Requires-Dist: tomli>=2.0; python_version < "3.11" and extra == "dev"

# 通过 Hologres DataFrame API 处理数据

- [功能介绍](#功能介绍)
- [安装](#安装)
- [创建 Session](#创建-session)
  - [连接池参数](#连接池参数)
  - [切换 Schema 和 Database](#切换-schema-和-database)
- [创建 DataFrame](#创建-dataframe)
  - [读取 Hologres 表和 SQL](#读取-hologres-表和-sql)
  - [从 Python 内存数据创建](#从-python-内存数据创建)
  - [从 pandas 创建](#从-pandas-创建)
  - [读取本地文件](#读取本地文件)
- [查询与转换](#查询与转换)
  - [列表达式](#列表达式)
  - [投影、过滤和排序](#投影过滤和排序)
  - [分组聚合](#分组聚合)
  - [Join 和 Union](#join-和-union)
  - [数组展开](#数组展开)
- [触发执行与获取结果](#触发执行与获取结果)
- [数据写入](#数据写入)
  - [追加写入](#追加写入)
  - [主键冲突处理](#主键冲突处理)
  - [覆盖写入](#覆盖写入)
- [OSS 文件读写](#oss-文件读写)
- [Serverless Computing](#serverless-computing)
- [Dynamic Table](#dynamic-table)
- [AI Function](#ai-function)
- [Function Compute UDF 和 UDTF](#function-compute-udf-和-udtf)
- [Catalog 和原生 SQL](#catalog-和原生-sql)
- [数据类型](#数据类型)
- [异常处理](#异常处理)
- [可运行示例](#可运行示例)
- [开发与测试](#开发与测试)

## 功能介绍

`hologres-dataframe` 是面向 Hologres 的 Python DataFrame API。接口采用 pandas 和主流
Python DataFrame API 的使用习惯，转换操作只构建惰性逻辑计划，在 action 触发时编译为
PostgreSQL 方言 SQL 并下推到 Hologres 执行。

主要能力：

- 使用 `select`、`filter`、`group_by`、`join` 等链式接口构建查询。
- 计算在 Hologres 内执行，只将最终结果返回 Python。
- 从 Hologres 表、SQL、Python 数据、pandas、本地文件和 OSS 文件创建 DataFrame。
- 将查询结果追加、UPSERT 或覆盖写入已有 Hologres 表。
- 支持 Dynamic Table、Serverless Computing、Hologres AI Function。
- 将 Python UDF/UDTF 自动部署为 Function Compute Remote UDX。
- 仅提供 snake_case 接口，例如 `group_by`，不提供 `groupBy` 别名。

```python
import hologres.dataframe as hg

with hg.connect(
    host="xxx.hologres.aliyuncs.com",
    port=80,
    dbname="my_db",
    user="<access_id>",
    password="<access_key>",
) as session:
    result = (
        session.table("public.orders")
        .filter(hg.col("amount") > 100)
        .with_column("tax", hg.col("amount") * 0.06)
        .select("order_id", "region", "amount", "tax")
        .sort(hg.col("amount"), ascending=False)
    )
    result.show(10)
```

`show()` 之前的操作都不会执行查询。`show()`、`collect()`、`count()`、`write()` 等
action 才会访问 Hologres。

## 安装

当前正式版本为 `0.1.0`：

```bash
pip install hologres-dataframe
```

安装指定版本：

```bash
pip install hologres-dataframe==0.1.0
```

从源码安装开发版本：

```bash
cd holo-dataframe-py
pip install -e ".[dev]"
```

依赖要求：

- Python >= 3.9
- `hologres-client>=0.1.2.dev0`
- `psycopg[binary,pool]>=3.1`
- `pandas>=1.3`
- `pyarrow>=12.0`
- `alibabacloud-fc20230330>=4.7.9,<5`

基础安装已经包含 `from_pandas()`、`to_pandas()` 和本地 CSV/JSON/Parquet 文件读取
所需依赖，也默认安装 Function Compute UDF/UDTF 自动部署 SDK。已有的
`hologres-dataframe[fc]` 安装命令仍然兼容，但不再需要额外指定 `[fc]`。

## 创建 Session

`Session` 管理连接池、数据库上下文、临时 Stage/表和 Function Compute 配置。建议作为
上下文管理器使用，退出时会释放连接并清理自动创建的临时对象。

```python
import hologres.dataframe as hg

with hg.connect(
    host="xxx.hologres.aliyuncs.com",
    port=80,
    dbname="my_db",
    user="<access_id>",
    password="<access_key>",
    sslmode="require",
) as session:
    session.table("public.orders").limit(10).show()
```

### 连接池参数

| 参数 | 默认值 | 说明 |
| --- | --- | --- |
| `min_size` | `1` | 连接池最小连接数 |
| `max_size` | `10` | 连接池最大连接数，也是并发上限 |
| `timeout` | `30.0` | 等待可用连接的最长秒数 |
| `max_idle` | `600.0` | 空闲连接回收时间 |
| `max_lifetime` | `3600.0` | 单个连接最大生命周期 |
| `application_name` | `hologres-dataframe` | `pg_stat_activity` 中显示的应用名 |
| `fc_config` | `None` | Function Compute UDF/UDTF 部署配置 |

其他关键字参数会透传给 psycopg/libpq，例如 `sslmode`、`connect_timeout`、
`keepalives` 和 `options`。

### 切换 Schema 和 Database

```python
session.use_schema("analytics")
print(session.current_schema())

session.use_database("another_db")
print(session.current_database())
```

`use_schema()` 修改当前 `search_path`。PostgreSQL 连接绑定 database，
`use_database()` 会关闭当前连接池并使用相同连接参数重建连接池。

## 创建 DataFrame

### 读取 Hologres 表和 SQL

```python
orders = session.table("public.orders")

paid_orders = session.sql(
    "SELECT order_id, user_id, amount FROM public.orders WHERE status = 'paid'"
)
```

`session.table()` 会读取并缓存表 Schema。`session.sql()` 返回惰性 DataFrame，原生 SQL
会在 action 时执行。

### 从 Python 内存数据创建

推荐从 Session 创建，使 DataFrame 可以直接执行 action 或写入：

```python
df = session.from_dict({"id": [1, 2], "name": ["alice", "bob"]})

df = session.from_records(
    [(1, "alice"), (2, "bob")],
    columns=["id", "name"],
)
```

小数据会编译为 SQL `VALUES`。超过 `inline_limit` 时会转为 Arrow 数据并上传到
Hologres。顶层 `hg.from_dict()`、`hg.from_records()` 适合仅构建 SQL 的场景；没有绑定
Session 的 DataFrame 不能执行 action。

### 从 pandas 创建

```python
df = session.from_pandas(pdf)
df.write("public.target")
```

数据在首次 action 时上传。Hologres 4.1 及以上版本使用 Internal Stage，旧版本使用临时
表兼容。Session 关闭时自动清理框架创建的 Stage 或临时表。

### 读取本地文件

```python
csv_df = session.read_files("data/*.csv")
parquet_df = session.read_files("data/orders.parquet")

csv_df.group_by("region").count().show()
```

本地文件支持 CSV、JSON、Parquet，通过 pyarrow 读取；文件读取发生在调用
`read_files()` 时，上传延迟到首次 action。

## 查询与转换

### 列表达式

`hg.col()` 引用列，`hg.lit()` 创建字面量。表达式中包含 `Column` 时只构建 SQL 表达式，
不会在 Python 中逐行计算。

```python
amount = hg.col("amount")

amount * 0.06
amount.between(100, 1000)
hg.col("region").isin("east", "west")
hg.col("name").ilike("ali%")
(amount > 100) & (hg.col("status") == "paid")
```

逻辑组合使用 `&`、`|`、`~`，每个比较条件需要加括号。Python 的 `and`、`or`、`not`
不能重载，不适用于 Column。

### 投影、过滤和排序

```python
result = (
    session.table("public.orders")
    .filter((hg.col("amount") > 100) & hg.col("region").is_not_null())
    .with_columns(
        {
            "tax": hg.col("amount") * 0.06,
            "amount_with_tax": hg.col("amount") * 1.06,
        }
    )
    .drop("internal_note")
    .rename({"region": "sales_region"})
    .select("order_id", "sales_region", "amount", "tax")
    .distinct()
    .order_by(hg.col("amount").desc())
    .limit(100)
    .offset(20)
)
```

`drop()` 和 `rename()` 只修改查询投影，不执行 `ALTER TABLE`。

### 分组聚合

```python
summary = (
    session.table("public.orders")
    .group_by("region")
    .agg(
        hg.sum("amount").alias("total_amount"),
        hg.avg("amount").alias("avg_amount"),
        hg.count("*").alias("order_count"),
    )
)
```

单一聚合可以使用快捷方法：

```python
session.table("public.orders").group_by("region").sum("amount")
session.table("public.orders").group_by("region").count()
```

长尾 SQL 函数使用 `hg.function()`，特殊 SQL 语法使用 `hg.expr()`：

```python
df.group_by("region").agg(
    hg.function("stddev", hg.col("amount")).alias("amount_stddev"),
    hg.expr("count(*) filter (where status = 'paid')").alias("paid_count"),
)
```

### Join 和 Union

```python
orders = session.table("public.orders")
users = session.table("public.users")

same_name_key = orders.join(users, on="user_id", how="left")

different_key = orders.join(
    users,
    on=orders["buyer_id"] == users["user_id"],
    how="inner",
)

all_rows = orders.union_all(session.table("public.orders_archive"))
deduplicated = orders.union(session.table("public.orders_archive"))
```

`on` 传字符串或字符串列表时编译为 `USING`；传 Column 表达式时编译为 `ON`。
`how` 支持 `inner`、`left`、`right`、`full`。

### 数组展开

```python
session.table("public.documents").explode("tags").show()
session.table("public.documents").explode("tags", outer=True).show()
```

`explode()` 编译为 Hologres 支持的 SELECT-list `unnest()`。`outer=True` 会保留数组为空
或 NULL 的原始行。

## 触发执行与获取结果

| Action | 返回值 | 适用场景 |
| --- | --- | --- |
| `show(n=10)` | `None` | 在终端查看少量结果 |
| `collect()` | `list[Row]` | 将完整结果加载到内存 |
| `to_pandas()` | `pandas.DataFrame` | 进入 pandas 处理流程 |
| `count()` | `int` | 服务端统计行数 |
| `first()` | `Row \| None` | 获取第一行 |
| `take(n)` | `list[Row]` | 获取前 n 行 |
| `iter_rows(batch_size=1000)` | Row 迭代器 | 分批读取大结果 |

```python
for row in session.table("public.orders").iter_rows(batch_size=2000):
    print(row.order_id, row.amount)
```

`collect()` 和 `to_pandas()` 会把完整结果放入客户端内存，大结果应优先使用
`iter_rows()`。

## 数据写入

目标表必须已经存在，建表和改表使用 `session.sql()`。`write()` 是立即执行的 action，
数据统一通过 `INSERT ... SELECT` 写入。

### 追加写入

```python
source = session.table("public.orders_staging").select("order_id", "user_id", "amount")
source.write("public.orders")
```

内存数据也使用同一接口：

```python
session.from_records(
    [(1, "east", 100), (2, "west", 200)],
    columns=["id", "region", "amount"],
).write("public.example_orders")
```

### 主键冲突处理

```python
# 主键冲突时忽略
df.write("public.target", on_conflict=("ignore", ["id"]))

# 只更新 DataFrame 中提供的非主键列
df.write("public.target", on_conflict=("update", ["id"]))

# 用 DataFrame 行替换目标行，未提供的目标列写 NULL
df.write("public.target", on_conflict=("replace", ["id"]))
```

`keys` 必须与目标表主键一致。

### 覆盖写入

```python
df.write("public.target", overwrite=True)

df.write(
    "public.partitioned_target",
    overwrite=True,
    partition={"ds": "2026-08-19"},
)
```

Hologres 3.1 及以上使用原生 `INSERT OVERWRITE`；低版本自动使用事务内
`TRUNCATE/DELETE + INSERT` 兼容方案。`overwrite=True` 与 `on_conflict` 互斥。

目标表启用 Binlog 时，overwrite 不能生成完整逐行 CDC 记录，默认会拒绝写入。确认下游允许
Binlog 缺口后，可以显式传入 `allow_binlog_gap=True`。

## OSS 文件读写

Hologres 4.1 及以上可以通过 `EXTERNAL_FILES` 读取 OSS 中的 CSV、Parquet 和 ORC：

```python
oss_df = session.read_files(
    "oss://bucket/input/",
    format="parquet",
    oss_endpoint="oss-cn-hangzhou-internal.aliyuncs.com",
    role_arn="acs:ram::123:role/hologres-oss-role",
)

oss_df.filter(hg.col("amount") > 100).show()
```

将查询结果导出到 OSS：

```python
session.table("public.orders").write_files(
    "oss://bucket/output/",
    format="csv",
    oss_endpoint="oss-cn-hangzhou-internal.aliyuncs.com",
    role_arn="acs:ram::123:role/hologres-oss-role",
    target_file_size_mb=128,
)
```

写 OSS 仅支持 CSV，路径必须是 `oss://`。`role_arn` 是 Hologres 访问 OSS 时扮演的 RAM
Role；已通过实例凭据授权时可以不传。导出本地文件可先调用 `to_pandas()`，再使用 pandas
的文件接口。

## Serverless Computing

```python
rows = (
    session.table("public.orders")
    .filter(hg.col("amount") > 100)
    .serverless(priority=4, required_cores=64, max_cores=128)
    .collect()
)
```

`serverless()` 是惰性标记，返回新的 DataFrame。执行 action 时，框架在同一事务和连接上
设置 `SET LOCAL`，不会把配置泄漏到连接池中的后续查询。

参数说明：

| 参数 | 范围 | 说明 |
| --- | --- | --- |
| `priority` | 1-5 | 查询优先级，5 最高 |
| `required_cores` | 0-102400 | 请求的 Core 数，0 表示自动估算 |
| `max_cores` | 0-102400 | 单查询 Core 上限，0 表示不设置查询级上限 |

Hologres 仍会判断 SQL 是否符合 Serverless 执行条件。不支持相关 GUC 的实例会直接报错，
客户端不会静默回退。

## Dynamic Table

```python
summary = (
    session.table("public.orders")
    .join(session.table("public.users"), on="user_id", how="left")
    .select("order_id", "user_id", "user_name")
)

summary.as_dynamic_table(
    "public.dt_orders",
    freshness="10 minutes",
    refresh_mode="incremental",
    auto_refresh=True,
    mode="ignore",
    options={"cdc_format": "binlog"},
)
```

`as_dynamic_table()` 是立即执行的 DDL action，返回 `None`。

| `mode` | 行为 |
| --- | --- |
| `ignore` | 默认；目标存在时不处理 |
| `errorifexists` | 目标存在时由 Hologres 报错 |
| `replace` | 在事务中删除并重建 |

手动刷新和查询：

```python
session.refresh_table("public.dt_orders")
session.table("public.dt_orders").show()
```

## AI Function

`hg.ai` 将 Hologres AI Function 暴露为惰性 Column 表达式：

```python
questions = session.table("public.questions")

questions.select(
    "question",
    hg.ai.gen(hg.function("concat", hg.lit("请简要回答："), hg.col("question"))).alias(
        "answer"
    ),
).show()

session.table("public.docs").select(
    "id",
    hg.ai.embed(hg.col("content")).alias("embedding"),
).show()

session.table("public.comments").select(
    "comment",
    hg.ai.sentiment(hg.col("comment")).alias("sentiment"),
).show()
```

还支持 `rank`、`chunk`、`classify`、`extract`、`mask`、`fix_grammar`、
`summarize`、`translate`、`similarity`、`parse_document`、`to_file` 和 `prompt`。
具体函数可用性取决于 Hologres 实例版本和模型配置。

## Function Compute UDF 和 UDTF

Hologres 没有内嵌 Python Runtime。`@hg.udf` 和 `@hg.udtf` 会在首次使用它们的 action
执行前，将 Python 代码部署到 Function Compute，再注册为 `LANGUAGE function_compute`
Remote UDX。

FC 配置放在 Session 中，AccessKey 使用阿里云标准凭据链，不要写在装饰器中：

```python
session = hg.connect(
    host="xxx.hologres.aliyuncs.com",
    port=80,
    dbname="my_db",
    user="<access_id>",
    password="<access_key>",
    fc_config={
        "endpoint": "123.cn-hangzhou-internal.fc.aliyuncs.com",
        "region": "cn-hangzhou",
        "role_arn": "acs:ram::123:role/fc-execution-role",
    },
)
```

标量 UDF：

```python
@hg.udf(packages=["numpy"])
def grade(score: int) -> str:
    return "high" if score >= 80 else "low"


session.table("public.scores").select("id", grade("score").alias("grade")).show()
```

UDTF 使用类和 `process()`：

```python
@hg.udtf(output_schema=hg.StructType([hg.StructField("word", hg.StringType())]))
class SplitWords:
    def process(self, text: str):
        for word in text.split():
            yield (word,)


session.table("public.docs").select("id", SplitWords("text")).show()
```

Hologres Remote UDX 当前不能返回 PostgreSQL `record`，因此 UDTF 的 `output_schema`
必须只有一个字段。FC 和 Hologres 必须在同一地域，并使用 FC 内网 endpoint。

## Catalog 和原生 SQL

```python
session.list_schemas()
session.list_tables()
session.list_tables("analytics")
session.current_schema()
session.current_database()
```

执行任意 Hologres SQL：

```python
session.sql("CREATE TABLE public.t (id bigint)").collect()
session.sql("INSERT INTO public.t VALUES (1)").collect()
```

`session.sql()` 本身是惰性的，即使传入 DDL/DML，也必须调用 action 才会执行。

## 数据类型

| Python API | Hologres/PostgreSQL 类型 |
| --- | --- |
| `BooleanType()` | `boolean` |
| `ByteType()` | `"char"` |
| `ShortType()` | `smallint` |
| `IntegerType()` | `integer` |
| `LongType()` | `bigint` |
| `FloatType()` | `real` |
| `DoubleType()` | `double precision` |
| `DecimalType(38, 2)` | `numeric(38,2)` |
| `StringType()` | `text` |
| `StringType(64)` | `varchar(64)` |
| `BinaryType()` | `bytea` |
| `DateType()` | `date` |
| `TimeType()` | `time` |
| `TimestampType()` | `timestamp` |
| `TimestampType(TimestampTimeZone.TZ)` | `timestamptz` |
| `JsonbType()` | `jsonb` |
| `GeographyType()` / `GeometryType()` | PostGIS 类型 |
| `ArrayType(StringType())` | `text[]` |
| `VectorType(float, 768)` | `vector(768)` |

`StructType` 和 `StructField` 用于描述 DataFrame 行结构、文件 Schema 和 UDTF 输出结构，
不是 Hologres 单列类型。

## 异常处理

所有公开异常都继承自 `HoloDataFrameError`：

```python
import hologres.dataframe as hg

try:
    session.table("public.orders").write(
        "public.target",
        overwrite=True,
        on_conflict=("update", ["id"]),
    )
except hg.InvalidArgumentError as exc:
    print(f"参数错误: {exc}")
except hg.UnsupportedOperationError as exc:
    print(f"当前 Hologres 版本不支持: {exc}")
except hg.ExecutionError as exc:
    print(f"Hologres 执行失败: {exc}")
```

| 异常 | 说明 |
| --- | --- |
| `InvalidArgumentError` | 参数类型、取值或接口组合不合法 |
| `CompilationError` | 逻辑计划无法编译为 SQL |
| `UnsupportedOperationError` | 当前引擎版本或能力不支持该操作 |
| `ExecutionError` | 连接、SQL 执行或服务端返回错误 |

## 可运行示例

[`examples`](examples/README.md) 目录提供直接连接 Hologres 的示例：

- [`basic_query.py`](examples/basic_query.py)：过滤、派生列、投影和排序。
- [`join_and_aggregate.py`](examples/join_and_aggregate.py)：Join 和分组聚合。
- [`write_rows.py`](examples/write_rows.py)：追加写入和 UPSERT。
- [`oss_files.py`](examples/oss_files.py)：OSS `EXTERNAL_FILES` 读写。
- [`dynamic_table.py`](examples/dynamic_table.py)：创建和刷新 Dynamic Table。

先配置连接环境变量：

```bash
export HOLO_HOST=xxx.hologres.aliyuncs.com
export HOLO_PORT=80
export HOLO_DATABASE=postgres
export HOLO_USER='<access-id>'
export HOLO_PASSWORD='<access-key>'

python examples/basic_query.py
```

各示例所需表结构和额外变量见 [`examples/README.md`](examples/README.md)。

## 开发与测试

```bash
cd holo-dataframe-py
python3 -m pytest
python3 -m ruff check hologres tests examples
python3 -m ruff format --check hologres tests examples
```

测试默认使用 fake connection 验证逻辑计划、SQL 编译和执行分发，不需要连接 Hologres。
需要真实实例的集成测试位于 `tests/integration`。

该发行包名为 `hologres-dataframe`，Python 导入路径为 `hologres.dataframe`。它与
`hologres-client` 共用 `hologres` namespace，并依赖后者提供高吞吐 Stage 写入等底层能力。
