You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何从Python(Pandas/Polars)高效写入Apache Iceberg?

Python环境下高效写入基于Hive Metastore的Apache Iceberg

我需要在Python环境中,通过Pandas(优先使用Polars)将数据高效写入基于Hive Metastore部署的Apache Iceberg,可选择直接写入或通过Trino查询引擎,目标是实现性能最优的写入方案。

目前已尝试多种方法,但写入吞吐量仍未达到预期,且相关指导文档较为匮乏,希望能找到高效写入Iceberg的正确方向(直接写入或通过Trino均可)。

已尝试的方法

  • 最初使用Trino Python客户端,将Pandas DataFrame解析为批量SQL查询:操作繁琐,速度最慢
  • 使用异步Trino Python包aiotrino,结合asyncio实现并发调用,利用Trino Worker资源:性能提升有限
  • 通过Trino的JDBC连接器(基于trino-jdbc.jar)结合jaydebeapi包:读取功能正常,但写入失败
  • 通过pySpark直接连接Iceberg:存在连接问题
  • pyIceberg 0.4.0版本:暂不支持记录写入功能(已查阅官方文档及包内代码确认)
  • 使用pandas.to_sql()方法:当前性能相对最优,但仍远未达标

另外了解到一种优化思路:先将DataFrame以Parquet格式直接上传至S3,再通过元数据调用关联Iceberg目录表(从压缩传输角度看具备合理性),但该方案尚未成功实现。

当前pandas.to_sql()实现代码

import warnings
# 忽略所有警告
warnings.filterwarnings("ignore")
import logging
logging.getLogger().setLevel(logging.DEBUG)

from sqlalchemy import create_engine
from trino.auth import BasicAuthentication
import pandas as pd 
from datetime import datetime

trino_target_host = "..." 
trino_target_port = 443
trino_target_catalog = 'iceberg'
trino_target_schema = '...'

def write_content(table: str, df, batch_size: int):
    print(f"size: {len(df)}")
    try:
        start = datetime.utcnow()
        engine = create_engine(
            f"trino://{trino_target_host}:{trino_target_port}/{trino_target_catalog}/{trino_target_schema}",
            connect_args={
                "auth": BasicAuthentication("...", "..."),
                "http_scheme": "https"
            }
        )
        # chunksize:分批次写入,节省内存
        # method='multi':一次插入多行,替代单条记录插入逻辑
        output = df.to_sql(table, engine, if_exists='append', index=False, chunksize=batch_size, method='multi')
        print(f"elapsed: {datetime.utcnow() - start}")
        return output
    except Exception as e:
        if 'Request Entity Too Large' in str(e):
            print(f"Batch size '{batch_size}' too large. Reduce size")
            return None
        print(f"failed: {e}")
        return None  

sample = df[:70000].copy(deep=True)
write_content('table_name', sample, 1000)

执行结果

size: 70000
EXECUTE IMMEDIATE not available for trino.dp.iotnxt.io:443; defaulting to legacy prepared statements (TrinoUserError(type=USER_ERROR, name=SYNTAX_ERROR, message="line 1:19: mismatched input ''SELECT 1''. Expecting: 'USING', ", query_id=20230629_161633_03350_rptf8))
elapsed: 0:14:05.315631

不同数据量下的读写性能对比(仅供参考,与部署环境强相关)

countReadwrite (batch = 1000)write (batch = 1500)
103.67 / 3.7 / 3.6 secs
1002.9 / 4.0 / 2.8 / 2.2 sec4.4 / 4.2 / 4.1 secs
10002.2 / 3.8 / 2.4 secs11.2 / 11.2 / 11.0 secs
10 0003.2 / 2.8 / 2.9 secs118.4 / 105.7 / 108.8 sec98.2 / 104.6 / 108.5 secs
70 0006.8 / 6.5 / 7.2 secs845.31 secs
140 0009.8 secs

内容的提问来源于stack exchange,提问作者Paul

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.17 09:52:42