从ClickHouse迁移到Doris集群,我是如何搞定Array字段的(附实战建表与Java代码)
从ClickHouse迁移到Doris集群:Array字段的实战迁移指南
当交通大数据平台的规模从几十个路口扩展到两百多个路口时,单点ClickHouse已经无法满足我们的性能需求。在迁移到Doris分布式集群的过程中,如何处理Array这样的复杂数据类型成为关键挑战。本文将分享我们在实际项目中从ClickHouse迁移到Doris时处理Array字段的完整经验。
1. 迁移背景与核心挑战
我们的平台需要处理多种交通数据:过车记录、违法抓拍、设备实时状态等。在小型项目中,ClickHouse的单机性能表现优异,但随着数据量和查询复杂度的增长,分布式架构成为必然选择。
迁移过程中,Array字段带来了三个主要挑战:
- 语法差异 :ClickHouse和Doris虽然都支持Array类型,但建表语法和函数接口存在细微差别
- 功能差异 :ClickHouse特有的arrayJoin函数在Doris中不存在,需要寻找替代方案
- 性能考量 :分布式环境下Array字段的存储和查询性能需要特别优化
提示:迁移前务必对两种数据库的Array支持进行详细对比,避免后期大规模返工
2. 数据类型对比与建表调整
2.1 Array类型支持对比
| 特性 | ClickHouse | Doris |
|---|---|---|
| 基础类型支持 | 全部标量类型 | 除JSON外的标量类型 |
| 嵌套数组支持 | 支持多层嵌套 | 仅支持单层 |
| 默认值设置 | 支持 | 有限支持 |
| 作为主键 | 不允许 | 不允许 |
| NULL处理 | 严格模式 | 宽松模式 |
2.2 建表语句迁移示例
原始ClickHouse建表语句:
CREATE TABLE traffic.security_metrics (
timestamp DateTime,
intersection_id UInt32,
safety_score Float32,
approach_metrics Array(Tuple(String, Float32, Float32))
) ENGINE = MergeTree()
ORDER BY (timestamp, intersection_id);
对应的Doris建表语句:
CREATE TABLE IF NOT EXISTS traffic.dwd_security_metrics (
`timestamp` DATETIME NOT NULL COMMENT '指标时间',
`intersection_id` INT NOT NULL COMMENT '路口ID',
`safety_score` FLOAT DEFAULT 0 COMMENT '安全评分',
`approach_metrics` ARRAY<VARCHAR(200)> COMMENT '进口指标数组'
)
DUPLICATE KEY(`timestamp`, `intersection_id`)
DISTRIBUTED BY HASH(`intersection_id`) BUCKETS 8
PROPERTIES (
"replication_num" = "3"
);
关键调整点:
- 移除了ClickHouse特有的ENGINE和ORDER BY语法
- 将Tuple类型转换为字符串拼接格式(Doris不支持嵌套复合类型)
- 明确指定了分布策略和副本数
3. 数据写入逻辑改造
3.1 Java数据拼接逻辑对比
ClickHouse原始代码:
List<Tuple3<String, Float, Float>> approaches = new ArrayList<>();
approaches.add(Tuple3.of("NB", 0.92f, 0.05f));
approaches.add(Tuple3.of("SB", 0.89f, 0.07f));
String sql = "INSERT INTO traffic.security_metrics VALUES (?, ?, ?, ?)";
PreparedStatement stmt = connection.prepareStatement(sql);
stmt.setObject(1, timestamp);
stmt.setObject(2, intersectionId);
stmt.setObject(3, safetyScore);
stmt.setObject(4, approaches);
Doris适配后的代码:
List<String> approaches = new ArrayList<>();
approaches.add("NB-0.92-0.05");
approaches.add("SB-0.89-0.07");
String arrayLiteral = "['" + String.join("','", approaches) + "']";
String sql = "INSERT INTO traffic.dwd_security_metrics VALUES (?, ?, ?, ?)";
PreparedStatement stmt = connection.prepareStatement(sql);
stmt.setObject(1, timestamp);
stmt.setObject(2, intersectionId);
stmt.setObject(3, safetyScore);
stmt.setObject(4, arrayLiteral);
3.2 批量写入优化
对于大规模数据迁移,建议使用Doris的Stream Load方式:
curl --location-trusted -u user:passwd \
-H "format: json" -H "strip_outer_array: true" \
-T data.json http://fe_host:8030/api/db/tbl/_stream_load
其中data.json格式示例:
[
{
"timestamp": "2023-05-01 08:00:00",
"intersection_id": 1001,
"safety_score": 0.85,
"approach_metrics": ["NB-0.92-0.05", "SB-0.89-0.07"]
}
]
4. 查询适配与性能优化
4.1 常用数组函数对照表
| ClickHouse函数 | Doris等效方案 | 示例 |
|---|---|---|
| arrayJoin | explode表函数 | 见4.2节 |
| arrayElement | element_at |
element_at(arr, 1)
|
| arrayMap | 组合使用transform等函数 |
transform(arr, x -> x*2)
|
| arrayFilter | 组合使用filter函数 |
filter(arr, x -> x > 0)
|
| length | array_size |
array_size(arr)
|
4.2 列转行实现方案
ClickHouse使用arrayJoin的典型查询:
SELECT
timestamp,
intersection_id,
arrayJoin(approach_metrics) AS metric
FROM traffic.security_metrics
在Doris中的等效实现:
SELECT
t.timestamp,
t.intersection_id,
e.item AS metric
FROM traffic.dwd_security_metrics t,
LATERAL explode(t.approach_metrics) e
对于复杂解析(如我们的字符串拼接格式):
SELECT
timestamp,
intersection_id,
split_part(trim(BOTH '"' FROM e.item), '-', 1) AS direction,
cast(split_part(trim(BOTH '"' FROM e.item), '-', 2) AS FLOAT) AS rate,
cast(split_part(trim(BOTH '"' FROM e.item), '-', 3) AS FLOAT) AS violation_rate
FROM traffic.dwd_security_metrics t,
LATERAL explode(t.approach_metrics) e
4.3 性能优化建议
- 合理设置分桶数 :Array字段经常参与查询时,应考虑按包含Array的字段分桶
- 避免过度嵌套 :Doris对多层嵌套支持有限,建议将复杂结构展平
- 使用物化视图 :对频繁查询的Array元素建立预计算
CREATE MATERIALIZED VIEW mv_approach_stats
DISTRIBUTED BY HASH(intersection_id)
REFRESH ASYNC
AS
SELECT
intersection_id,
element_at(metrics, 1) AS direction,
element_at(metrics, 2) AS rate
FROM (
SELECT
intersection_id,
explode(approach_metrics) AS metrics
FROM dwd_security_metrics
) t
5. 迁移后的验证与监控
完成迁移后,我们建立了完整的验证流程:
-
数据一致性检查 :编写对比脚本验证两边数据是否一致
# 抽样对比脚本示例 def compare_array(ch_data, doris_data): ch_set = set(tuple(x) for x in ch_data) doris_set = set(tuple(x.split('-')) for x in doris_data) return ch_set == doris_set -
查询性能基准测试 :对典型查询场景进行性能对比
查询类型 ClickHouse(ms) Doris(ms) 差异分析 单路口点查 12 18 网络开销增加 全区域聚合 320 210 分布式优势显现 数组展开分析 45 60 函数效率差异 -
资源监控配置 :特别关注Array字段相关的指标
- BE节点内存使用
- 数组函数CPU消耗
- 网络传输量
在实际运行中,我们发现Doris的分布式特性确实很好地支撑了大规模数据的处理,虽然某些数组操作性能略低于ClickHouse,但通过合理的分桶策略和查询优化,最终用户体验差异不大。
更多推荐
所有评论(0)