如何从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
不同数据量下的读写性能对比(仅供参考,与部署环境强相关)
| count | Read | write (batch = 1000) | write (batch = 1500) |
|---|---|---|---|
| 10 | 3.67 / 3.7 / 3.6 secs | ||
| 100 | 2.9 / 4.0 / 2.8 / 2.2 sec | 4.4 / 4.2 / 4.1 secs | |
| 1000 | 2.2 / 3.8 / 2.4 secs | 11.2 / 11.2 / 11.0 secs | |
| 10 000 | 3.2 / 2.8 / 2.9 secs | 118.4 / 105.7 / 108.8 sec | 98.2 / 104.6 / 108.5 secs |
| 70 000 | 6.8 / 6.5 / 7.2 secs | 845.31 secs | |
| 140 000 | 9.8 secs |
内容的提问来源于stack exchange,提问作者Paul
相关产品推荐
相关产品推荐

