过去,Doris 与 Iceberg 的组合更像一条单行道:Iceberg 负责保存湖上数据,Doris 负责高速查询。Apache Doris 4.1 补上的正是写入和治理能力,包括 UPDATE、DELETE、MERGE INTO,表结构与分区演进,以及 rewrite_data_files、expire_snapshots 等维护操作,同时完整支持 Iceberg V3。变化的重点不只是 SQL 变多了,而是团队可以在同一个分析入口内完成查询、修正、合并和维护。
从“能查”走向“能管理”
只支持查询时,数据修正通常需要回到 Spark、Flink 或其他计算引擎。一次看似简单的状态更新,可能意味着额外编排任务、重复配置 Catalog,并处理不同引擎之间的权限和语义差异。
Doris 4.1 支持数据修改后,典型工作流可以缩短为:
- 用
UPDATE修正少量已知记录。 - 用
DELETE清理过期、错误或不再合规的数据。 - 用
MERGE INTO将增量数据按业务键合并到主表。 - 在同一套湖表上继续执行交互式分析。
其中 MERGE INTO 尤其关键。增量同步经常同时包含新增与变更,单独执行插入和更新不仅步骤多,还容易在重试时制造重复数据。将匹配、更新和插入规则写进一条语句,更适合 CDC 落湖、维表同步和批次补数。
一套可以改造的增量合并示例
下面示例假设 Doris 已配置可写的 Iceberg Catalog,并且目标表与暂存表位于该 Catalog 中。Catalog 创建参数、对象存储地址和认证方式需要按实际环境替换;示例用于展示 Doris 4.1 中可以这样组织湖表写入流程。
-- 目标表:保存当前订单状态
CREATE TABLE iceberg_catalog.sales.orders (
order_id BIGINT,
customer_id BIGINT,
order_status STRING,
amount DECIMAL(18, 2),
updated_at TIMESTAMP
)
PARTITION BY (days(updated_at));
-- 暂存表:接收本批增量
CREATE TABLE iceberg_catalog.sales.orders_delta (
order_id BIGINT,
customer_id BIGINT,
order_status STRING,
amount DECIMAL(18, 2),
updated_at TIMESTAMP
);
MERGE INTO iceberg_catalog.sales.orders AS target
USING iceberg_catalog.sales.orders_delta AS source
ON target.order_id = source.order_id
WHEN MATCHED AND source.updated_at >= target.updated_at THEN
UPDATE SET
customer_id = source.customer_id,
order_status = source.order_status,
amount = source.amount,
updated_at = source.updated_at
WHEN NOT MATCHED THEN
INSERT (order_id, customer_id, order_status, amount, updated_at)
VALUES (source.order_id, source.customer_id, source.order_status,
source.amount, source.updated_at);
这里用 updated_at 限制旧批次覆盖新数据,是生产环境中值得保留的保护条件。实际接入 CDC 时,还应明确删除事件如何表达、业务键是否稳定,以及任务重试是否会重复提交。具体 SQL 语法应以部署版本和 Catalog 配置为准。
对于范围明确的数据修正,则可以直接使用更短的语句:
UPDATE iceberg_catalog.sales.orders
SET order_status = 'CANCELLED',
updated_at = CURRENT_TIMESTAMP
WHERE order_id = 10001
AND order_status = 'PENDING';
DELETE FROM iceberg_catalog.sales.orders
WHERE updated_at < TIMESTAMP '2023-01-01 00:00:00'
AND order_status = 'CANCELLED';
执行大范围更新或删除前,建议先用同样的 WHERE 条件运行 SELECT COUNT(*),确认影响行数,并保留可回退的 Iceberg 快照。
表结构和分区不再一次定死
数据湖表会长期存在,业务字段和查询方式却一直变化。完整的表结构管理与分区演进意味着团队不必因为新增字段或分区策略调整就整体重写表。
例如,订单表上线后可能需要增加渠道字段;随着数据规模增长,原有分区粒度也可能不再合适。可以这样实践变更流程:
ALTER TABLE iceberg_catalog.sales.orders
ADD COLUMN sales_channel STRING;
-- 分区演进的具体表达方式依 Doris 4.1 与 Iceberg Catalog 的语法而定。
-- 上线前先在测试表验证新旧分区能否被统一读取和正确裁剪。
分区演进并不会自动改写历史文件。新写入数据采用新规则,旧数据仍保留原布局,查询引擎需要同时理解两代分区规范。这也是 Iceberg 的价值之一:把这种演进记录在表元数据中,而不是要求应用拼接不同目录。
写得越频繁,维护越不能缺席
更新、删除和小批次合并会持续产生数据文件、删除文件、快照与元数据。写入链路跑通并不代表表会长期保持健康。Doris 4.1 支持 rewrite_data_files 和 expire_snapshots 等操作,让日常维护也能纳入统一工作流。
rewrite_data_files用于合并碎片化的小文件,降低文件枚举、打开和扫描开销。expire_snapshots用于清理超过保留策略的历史快照及相关文件,控制存储和元数据增长。
维护任务不应盲目高频执行。文件重写会消耗计算和 I/O,快照过早过期则会缩短审计、回滚和故障恢复窗口。比较稳妥的做法是先监控文件数量、平均文件大小、快照数量和查询延迟,再设置触发阈值。维护过程的具体调用语法和参数应按 Doris 4.1 的过程接口及所用 Catalog 进行配置。
Iceberg V3 带来的采用边界
完整支持 Iceberg V3,使 Doris 可以参与采用新版表格式的湖仓架构。但“Doris 能读写”不等于整条数据链路都已兼容。Spark、Flink、采集工具、Catalog 服务和治理平台都可能访问同一张表,任何一个旧组件都可能成为升级阻点。
正式迁移前,至少检查以下事项:
- 盘点所有读写 Iceberg 表的引擎及其 V3 支持情况。
- 用隔离表验证
UPDATE、DELETE、MERGE INTO的结果和并发行为。 - 验证表结构演进后,旧任务与新查询是否都能正确解析字段。
- 为快照保留期、文件重写频率和失败重试建立明确策略。
- 先迁移可回放、低风险的数据集,再处理核心事实表。
Doris 4.1 的意义,在于把 Iceberg 从外部只读数据源推进为可操作的数据资产。真正落地时,不要只验证一条 MERGE INTO 能否执行;还要把多引擎兼容、并发写入、快照保留和文件维护一起纳入验收。