
# DataFrame示例
本章节提供了类Pandas的Python DataFrame SDK，方便用户使用Python编写数据处理作业；同时利用Aura内核高效的计算能力，为数据科学家、AI工程师等提供了易用、高效的数据处理能力。
本特性基于Ibis Python DataFrame开源框架实现，将Ibis前端框架与Aura引擎对接。用户可基于熟悉的Ibis DataFrame API编写数据处理脚本，Ibis框架将Python API翻译为Aura引擎可执行的SQL语句并下发，从而实现计算逻辑在Aura引擎的高效处理。支持ibis的api使用，具体ibis api请参见[Ibis官方文档](https://ibis-project.org/)。
#### 快速开始
以下代码使用ibis库连接Aura并执行数据查询，将结果转换为DataFrame格式的基本语法。
示例仅供参考，请您根据实际情况进行修改。
关于Ibis更详细的用法，请参见[Ibis官方文档](https://ibis-project.org/)。
```
import ibis  # 导入ibis依赖
import logging
import os
from aura_frame.multimodal import ai_lake
con = ai_lake.connect(
    aura_endpoint=os.getenv("aura_endpoint"), # 指定服务的区域，区域查询请参见地区和终端节点
    aura_endpoint_name=os.getenv("aura_endpoint_name"),
    aura_workspace_id=os.getenv("aura_workspace_id"), # 获取workspace_id
    lf_catalog_name=os.getenv("lf_catalog_name"), # 连接指定的Catalog
    access_key=os.getenv("access_key"),  # 获取ak/sk
    secret_key=os.getenv("secret_key"), 
    use_single_cn_mode=True,  # 表示是否开启单CN模式
    logging_level=logging.INFO, # 设置日志级别
)
con.set_function_staging_workspace(
    obs_directory_base=os.getenv("obs_directory_base"), # obs中udf的存储路径
    obs_bucket_name=os.getenv("obs_bucket_name"), # obs的名称
    obs_server=os.getenv("obs_server"), # obs访问地址,详情请参见终端节点（Endpoint）和访问域名
    access_key=os.getenv("access_key"),  
    secret_key=os.getenv("secret_key")
)
ds = con.load_dataset("table_name", database="db")  # 通过连接到后端获取table表信息，建立表对象
ds.select_columns([ds.cola])
df = ds.execute()  # 将DataFrame转为SQL，传输到后端执行，并且返回Pandas DataFrame格式的结果
con.close()
```
#### 不带UDF的DF示例
下文以tpch的query1为例，展示DataFrame的用法。
查询SQL为：
```
SELECT
    l_returnflag,
    l_linestatus,
    sum(l_quantity) AS sum_qty,
    sum(l_extendedprice) AS sum_base_price,
    sum(l_extendedprice * (1 - l_discount)) AS sum_disc_price,
    sum(l_extendedprice * (1 - l_discount) * (1 + l_tax)) AS sum_charge,
    avg(l_quantity) AS avg_qty,
    avg(l_extendedprice) AS avg_price,
    avg(l_discount) AS avg_disc,
    count(*) AS count_order
FROM
    lineitem
WHERE
    l_shipdate <= CAST('1998-09-02' AS date)
GROUP BY
    l_returnflag,
    l_linestatus
ORDER BY
    l_returnflag,    l_linestatus;
```
对应的DataFrame逻辑如下：
```
import os
from datetime import datetime, timedelta
from aura_frame.multimodal import ai_lake
target_database = "test"
con = ai_lake.connect(
    aura_endpoint=os.getenv("aura_endpoint"),
    aura_endpoint_name=os.getenv("aura_endpoint_name"),
    aura_workspace_id=os.getenv("aura_workspace_id"),
    lf_catalog_name=os.getenv("lf_catalog_name"),
    access_key=os.getenv("access_key"),
    secret_key=os.getenv("secret_key"),
    default_database=target_database,
    use_single_cn_mode=True,
)
try:
    ds = con.load_dataset("lineitem", database="tpch")
    cutoff = datetime(1998, 12, 1) - timedelta(days=90)
    q = ds.filter(expr=ds.l_shipdate <= cutoff)
    discount_price = ds.l_extendedprice * (1 - ds.l_discount)
    charge = discount_price * (1 + ds.l_tax)
    q = q.groupby(["l_returnflag", "l_linestatus"])
    q = q.aggregate(
        sum_qty=ds.l_quantity.sum(),
        sum_base_price=ds.l_extendedprice.sum(),
        sum_disc_price=discount_price.sum(),
        sum_charge=charge.sum(),
        avg_qty=ds.l_quantity.mean(),
        avg_price=ds.l_extendedprice.mean(),
        avg_disc=ds.l_discount.mean(),
        count_order=ds.count(),
    )
    q = q.order_by(["l_returnflag", "l_linestatus"])
    df = q.execute()
finally:
    con.close()
```
#### 带Scalar UDF的DF示例
**场景描述**
在AI数据工程中，数据预处理是一个关键步骤，通常需要对存储在数据库中的数据进行复杂的清洗、转换和特征工程操作。然而，传统的数据预处理逻辑往往在数据库外部通过Python脚本实现，这会导致大量数据在数据库和Python环境之间传输，不仅增加了计算开销，还无法充分利用数据库的分布式计算能力。通过在数据库内核中实现Python UDF，并在aura_frame中增加Python UDF显式，隐式调用接口，可以将数据预处理逻辑直接嵌入到数据库查询中，减少数据传输，提升整体处理效率。
对于复杂业务逻辑，例如使用机器学习模型对特征数据表进行按行过滤。传统函数式UDF定义方式无法有效支持模型初始化和资源管理需求，因为模型超参数和资源（如模型文件加载）需要在函数执行前一次性配置，而函数式UDF无法高效管理此类初始化逻辑和资源释放操作。通过以Class形式定义Python UDF，用户可以在__init__方法中配置模型超参数和初始化资源（如加载模型），并在__call__方法中执行推理逻辑，同时通过__del__方法释放资源（如关闭文件或连接），从而实现高效的资源管理和业务逻辑处理。
**约束限制**
- Python UDF功能约束限制如下：
  - 仅支持Python语言编写的用户自定义函数。
  
  - Python环境依赖管理由用户自行负责，需确保各依赖库版本兼容。
  
  - 运行时环境要求Python版本为3.11。
   
- Python Class UDF功能约束限制如下：
  - 类必须定义__init__方法，且其参数名称需与调用时的参数名称保持一致。
  
  - 类必须定义__call__成员方法，该方法为UDF的主函数入口。
  
  - 类可按需定义__del__成员方法，且该方法不支持传入参数。
  
  - 调用class型UDF时，__init__方法的参数仅支持传入常量，不支持传入表的数据列或其他表达式。
   
结合DataFrame使用Scalar UDF是推荐的标准用法，此时整个Scalar UDF的外部必须要包围DataFrame的SELECT方法。
经过注册后返回的值是DataFrame中的一个UDF算子。此时，该算子可以被多个DataFrame表达式多次调用，示例如下：
```
import os
import aura_frame as aura
from aura_frame.multimodal import ai_lake
def transform_json(ts: float, msg: str) -> str:
    import json
    from car_bu.parser_core import ParseObjects, dict_to_object
    if msg == '0_msg':
        return json.dumps({"time_stamp": 0.0, "msg": {}})
    else:
        d = dict_to_object(json.loads(msg))
        return json.dumps({"time_stamp": ts / 10, "msg": ParseObjects(d)})
target_database = "test"
con = ai_lake.connect(
    aura_endpoint=os.getenv("aura_endpoint"),
    aura_endpoint_name=os.getenv("aura_endpoint_name"),
    aura_workspace_id=os.getenv("aura_workspace_id"),
    lf_catalog_name=os.getenv("lf_catalog_name"),
    access_key=os.getenv("access_key"),
    secret_key=os.getenv("secret_key"),
    default_database=target_database,
    use_single_cn_mode=True,
)
try:
    udf = con.create_scalar_function(
        transform_json,
        database=target_database,
        imports=["car_bu/parser_core.py"],
    )
    ds = con.load_dataset("your-table", database=target_database)
    ds = ds.map(fn=udf, on=[ds.ts, ds.msg], as_col="json_column")
    ds = ds.select_columns(ds.ts, ds.msg, ds.json_column)
    df = ds.execute()
finally:
    con.close()
```
对于Python类的Scalar Class UDF，使用方式和Python函数的Scalar UDF稍有不同：UDF算子接受__call__方法参数；with_arguments接受__init__方法参数。完整示例如下：
```
import os
from aura_frame.multimodal import ai_lake
from sentencepiece import SentencePieceProcessor
target_database = "test"
class SPManager:
    def __init__(self, model_file: str, *, bos: bool = False, eos: bool = False, reverse: bool = False) -> None:
        sp = SentencePieceProcessor(model_file=model_file)
        opts = ":".join(opt for opt, flag in (("bos", bos), ("eos", eos), ("reverse", reverse)) if flag)
        sp.set_encode_extra_options(opts)
        self.spm = sp
    def __call__(self, input_text: str) -> str:
        import json
        return json.dumps(self.spm.encode(input_text, out_type=str))
con = ai_lake.connect(
    aura_endpoint=os.getenv("aura_endpoint"),
    aura_endpoint_name=os.getenv("aura_endpoint_name"),
    aura_workspace_id=os.getenv("aura_workspace_id"),
    lf_catalog_name=os.getenv("lf_catalog_name"),
    access_key=os.getenv("access_key"),
    secret_key=os.getenv("secret_key"),
    default_database=target_database,
    use_single_cn_mode=True,
)
try:
    con.create_scalar_function(
        SPManager,
        database=target_database,
        packages=["sentencepiece==0.2.0"],
        imports=["test_model.model"],
    )
    udf = con.get_function("SPManager", database=target_database)
    ds = con.load_dataset("your-table", database=target_database)
    ds = ds.map(
        fn=udf, on=[ds.data], as_col="pieces_column",
        model_file="test_model.model", bos=True, eos=True,
    )
    df = ds.execute()
finally:
    con.close()
```
#### 带UDAF的DF示例
结合DataFrame使用UDAF是推荐的标准用法，此时整个UDAF的外部必须要包围DataFrame的aggregate方法。
UDAF只能出现在SELECT、ORDER BY、HAVING 3种从句中。
经过注册后返回的值是DataFrame中的一个UDF算子。此时，该算子可以被多个DataFrame表达式多次调用，with_arguments接受__init__方法参数，示例如下：
```
import os
import json
from aura_frame.multimodal import ai_lake
from aura_frame.multimodal.function import AggregateFnBuilder
target_database = "test"
class PythonTopK:
    def __init__(self, k: int):
        self.k = k
        self._agg_state = []
    @property
    def aggregate_state(self):
        return self._agg_state
    def accumulate(self, value: int) -> None:
        if value is not None:
            self._agg_state.append(value)
            self._agg_state.sort(reverse=True)
            self._agg_state = self._agg_state[: self.k]
    def merge(self, other_state) -> None:
        merged = self._agg_state + other_state
        merged.sort(reverse=True)
        self._agg_state = merged[: self.k]
    def finish(self) -> str:
        return json.dumps(self._agg_state)
con = ai_lake.connect(
    aura_endpoint=os.getenv("aura_endpoint"),
    aura_endpoint_name=os.getenv("aura_endpoint_name"),
    aura_workspace_id=os.getenv("aura_workspace_id"),
    lf_catalog_name=os.getenv("lf_catalog_name"),
    access_key=os.getenv("access_key"),
    secret_key=os.getenv("secret_key"),
    default_database=target_database,
    use_single_cn_mode=True,
)
try:
    udaf = con.create_agg_function(PythonTopK, database=target_database)
    ds = con.load_dataset("your-table", database=target_database)
    # 1. SELECT
    udaf_builder1 = AggregateFnBuilder(fn=udaf, on=[ds.amount], as_col="top3", k=3)
    summary = ds.aggregate([udaf_builder1], by=[ds.category])
    df = summary.execute()
    # 2. ORDER BY
    udaf_builder2 = AggregateFnBuilder(fn=udaf, on=[ds.amount], as_col="sort", k=10)
    summary = ds.aggregate([udaf_builder2], by=[ds.category])
    summary = summary.order_by(ds.amount.desc())
    df = summary.execute()
    # 3. HAVING
    udaf_builder3 = AggregateFnBuilder(fn=udaf, on=[ds.amount], as_col="top", k=5)
    summary = ds.aggregate([udaf_builder3], by=[ds.category])
    df = summary.execute()
finally:
    con.close()
```
#### 直接使用带Scalar UDF的DF示例
**场景描述**
在大数据处理场景中，用户在使用DataFrame进行数据处理时，常常需要通过用户自定义函数（UDF）来实现复杂的数据计算逻辑。然而，当前系统中UDF的注册与调用耦合在一起，用户无法在注册后单独查看或删除已注册的UDF，这在团队协作开发或动态管理UDF时带来了诸多不便。为了解决这一问题，本次需求通过新增Backend.udf系列API，支持用户在运行时动态查看、调用和删除UDF，从而提升UDF管理的灵活性和开发效率。
**约束限制**
UDF直接调用/查看/删除功能约束限制如下：
用户必须先获得Backend（Aura）连接，才能调用Backend UDF Registry API。
对于复杂类型的支持依赖Aura内核对复杂类型的支持。
在注册和使用分离的场景下，对于UDF的使用者，允许用户通过Aura Backend提供的Backend.udf接口直接获得、查看、删除Scalar UDF，示例如下：
```
import os
from aura_frame.multimodal import ai_lake
target_database = "test"
con = ai_lake.connect(
    aura_endpoint=os.getenv("aura_endpoint"),
    aura_endpoint_name=os.getenv("aura_endpoint_name"),
    aura_workspace_id=os.getenv("aura_workspace_id"),
    lf_catalog_name=os.getenv("lf_catalog_name"),
    access_key=os.getenv("access_key"),
    secret_key=os.getenv("secret_key"),
    default_database=target_database,
    use_single_cn_mode=True,
)
try:
    udf_names = con._conn.udf.names(database=target_database)
    if "transform_json" in udf_names:
        transform_json_udf = con.get_function("transform_json", database=target_database)
        t = con.load_dataset("your-table", database=target_database)
        expression = t.select(transform_json_udf(t.ts, t.msg).name("json_column"))
        df = expression.execute()
        con.delete_function("transform_json", database=target_database, if_it_exists=True)
    if "SPManager" in udf_names:
        sp_udf = con.get_function("SPManager", database=target_database)
        t = con.load_dataset("your-table", database=target_database)
        expression = t.select(
            sp_udf(t.data)
            .with_arguments(model_file="test_model.model", bos=True, eos=True)
            .name("pieces_column")
        )
        df = expression.execute()
        con.delete_function("SPManager", database=target_database, if_it_exists=True)
finally:
    con.close()
```
#### 直接使用带UDAF的DF示例
在注册和使用分离的场景下，对于UDAF的使用者，允许用户通过Aura Backend提供的Backend.udaf接口直接获得、查看、删除UDAF，示例如下：
```
import os
from aura_frame.multimodal import ai_lake
target_database = "test"
con = ai_lake.connect(
    aura_endpoint=os.getenv("aura_endpoint"),
    aura_endpoint_name=os.getenv("aura_endpoint_name"),
    aura_workspace_id=os.getenv("aura_workspace_id"),
    lf_catalog_name=os.getenv("lf_catalog_name"),
    access_key=os.getenv("access_key"),
    secret_key=os.getenv("secret_key"),
    default_database=target_database,
    use_single_cn_mode=True,
)
try:
    udaf_names = con._conn.udaf.names(database=target_database)
    if "pythontopk" in udaf_names:
        topk_udaf = con._conn.udaf.get(name="pythontopk", database=target_database)
        t = con.load_dataset("your-table", database=target_database)
        expression = t.aggregate(
            top_values=topk_udaf(t.amount).with_arguments(k=3),
            by=[t.category],
        )
        df = expression.execute()
        con._conn.udaf.unregister("pythontopk", database=target_database)
finally:
    con.close()
```
#### DataFrame create_table函数使用说明
create_table用于在Aura创建一个表，函数签名如下：
```
def create_table(
    self,
    name: str,
    obj: Optional[Union[ir.Table, pd.DataFrame, pa.Table, pl.DataFrame, pl.LazyFrame]] = None,
    *,
    schema: Optional[ibis.Schema] = None,
    database: Optional[str] = None,
    temp: bool = False,
    external: bool = False,
    overwrite: bool = False,
    partition_by: Optional[ibis.Schema] = None,
    table_properties: Optional[dict] = None,
    store: Optional[str] = None,
    location: Optional[str] = None
)
```
表1参数说明 
| 参数名称             | 类型                                                           | 是否必须 | 说明                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                           |
|:---|:---|:---|:---|
| name             | str                                                          | 是    | 要创建的表名。                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                      |
| obj              | ir.Table\|pd.DataFrame\|pa.Table\|pl.DataFrame\|pl.LazyFrame | 否    | 用于填充表格的数据；必须至少指定obj或schema之一（当前不支持在创建表的时候插入数据）。                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                              |
| schema           | sch.SchemaLike                                               | 否    | 要创建的表的架构；必须至少指定obj或schema之一。                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                 |
| database         | str                                                          | 否    | 创建表的数据库的名称；如果未传递，则使用当前数据库。                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                   |
| temp             | bool                                                         | 否    | 是否创建为临时表。默认为False。                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                           |
| external         | bool                                                         | 否    | 是否创建为外表。默认为False。                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                            |
| overwrite        | bool                                                         | 否    | 如果为True，表已存在则替换表。默认为False（当前不支持覆盖表）。                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                         |
| partition_by     | sch.SchemaLike                                               | 否    | 指定分区列，分区列中出现的列不能出现在表的普通列描述中。                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                 |
| table_properties | dict                                                         | 否    | 表级别可选参数设置，支持参数范围如下参考[表1](https://support.huaweicloud.com/sqlref-aura-aidatalake/aidatalake_061_0219.html#aidatalake_061_0219__zh-cn_topic_0000002560961201_table455155917451)。                                                                                                                                                                                                                                                                                                                                                                                               |
| store            | str                                                          | 否    | 表存储格式，支持ORC、PARQUET、HUDI、ICEBERG四种存储格式。                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                      |
| location         | str                                                          | 否    | 表存储路径，必须为合法OBS路径，支持OBS对象桶和并行文件系统。 - 如果该路径为OBS对象桶路径，则该表只读，否则该表支持读写。  - 如果创建表类型为托管表（Managed Table），则不允许指定表存储路径，表存储路径由系统指定为Schema路径下与表名同名路径，且要求该路径在建表时为空。   |
   
