ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

【中台·数据篇】存储与计算:Hive、Spark 与 Flink 引擎选型与架构

【中台·数据篇】存储与计算:Hive、Spark 与 Flink 引擎选型与架构 前言数据采集到 Kafka 后需要存储和计算引擎处理。本篇详解 Hive数仓存储、Spark离线计算、Flink实时计算三大引擎的选型和架构设计。一、存储引擎对比┌──────────────────────────────────────────────────────────┐ │ 数据存储引擎选择 │ │ │ │ ┌──────────┐ ┌──────────┐ ┌──────────┐ ┌──────────┐│ │ │ HDFS │ │ Hive │ │ClickHouse│ │ HBase ││ │ │ 文件存储 │ │ 数仓SQL │ │ 实时OLAP │ │ KV存储 ││ │ └──────────┘ └──────────┘ └──────────┘ └──────────┘│ │ ┌──────────┐ ┌──────────┐ │ │ │ Redis │ │ ES │ │ │ │ 缓存 │ │ 全文检索 │ │ │ └──────────┘ └──────────┘ │ └──────────────────────────────────────────────────────────┘存储用途特点查询延迟HDFS原始数据/归档海量、便宜、不可变秒-分钟Hive数仓 SQL 查询SQL on HDFS分钟ClickHouse实时分析列存、向量化毫秒-秒HBase高频点查行存、低延迟毫秒Redis实时指标内存、极快微秒ES全文检索倒排索引毫秒二、Hive 数仓架构SQL 查询 → Hive Server2 → 解析器 → 优化器 → 执行器 ↓ MapReduce / Tez / Spark ↓ HDFS数据文件建表-- ODS 层原始数据 CREATE TABLE IF NOT EXISTS ods_orders ( order_id BIGINT, user_id BIGINT, amount DECIMAL(10,2), status STRING, create_time TIMESTAMP ) PARTITIONED BY (dt STRING) -- 按天分区 STORED AS ORC -- ORC 列存格式 LOCATION /data/warehouse/ods/orders; -- DWD 层清洗明细 CREATE TABLE IF NOT EXISTS dwd_orders ( order_id BIGINT, user_id BIGINT, amount DECIMAL(10,2), status STRING, payment_method STRING, create_time TIMESTAMP, province STRING, city STRING ) PARTITIONED BY (dt STRING) STORED AS ORC LOCATION /data/warehouse/dwd/orders; -- DWS 层用户主题汇总 CREATE TABLE IF NOT EXISTS dws_user_daily ( user_id BIGINT, order_count INT, total_amount DECIMAL(12,2), avg_amount DECIMAL(10,2), last_order_time TIMESTAMP ) PARTITIONED BY (dt STRING) STORED AS ORC LOCATION /data/warehouse/dws/user_daily; -- ADS 层应用指标 CREATE TABLE IF NOT EXISTS ads_daily_sales ( dt STRING, total_orders BIGINT, total_amount DECIMAL(15,2), avg_amount DECIMAL(10,2), paying_users BIGINT, new_users BIGINT ) STORED AS ORC LOCATION /data/warehouse/ads/daily_sales;ETL 示例-- ODS → DWD清洗 INSERT OVERWRITE TABLE dwd_orders PARTITION (dt2024-01-15) SELECT order_id, user_id, CAST(amount AS DECIMAL(10,2)) AS amount, UPPER(status) AS status, payment_method, create_time, -- 从 IP 推导省市 ip_to_province(ip) AS province, ip_to_city(ip) AS city FROM ods_orders WHERE dt 2024-01-15 AND order_id IS NOT NULL; -- 过滤脏数据 -- DWD → DWS用户日汇总 INSERT OVERWRITE TABLE dws_user_daily PARTITION (dt2024-01-15) SELECT user_id, COUNT(*) AS order_count, SUM(amount) AS total_amount, AVG(amount) AS avg_amount, MAX(create_time) AS last_order_time FROM dwd_orders WHERE dt 2024-01-15 AND status PAID GROUP BY user_id; -- DWS → ADS日销售指标 INSERT OVERWRITE TABLE ads_daily_sales SELECT 2024-01-15 AS dt, COUNT(DISTINCT user_id) AS paying_users, SUM(order_count) AS total_orders, SUM(total_amount) AS total_amount, AVG(total_amount / order_count) AS avg_amount FROM dws_user_daily WHERE dt 2024-01-15;分区与分桶-- 分区按时间天 PARTITIONED BY (dt STRING) -- 查询时只扫需要的分区 SELECT * FROM dwd_orders WHERE dt 2024-01-15; -- 分桶按字段如 user_id CLUSTERED BY (user_id) INTO 32 BUCKETS -- 提升采样查询和 JOIN 性能存储格式-- ORC推荐列存、压缩、索引 STORED AS ORC TBLPROPERTIES ( orc.compress SNAPPY, orc.create.index true ); -- Parquet列存、跨平台 STORED AS PARQUET; -- Text行存、不压缩ODS 临时用 STORED AS TEXTFILE;三、Spark 计算引擎架构Spark Driver ├── SparkContext ├── DAG Scheduler → 任务拆分 └── Task Scheduler → 任务分发 Cluster ManagerYARN/K8s ├── Worker 1 → Executor多 Task ├── Worker 2 → Executor └── Worker 3 → ExecutorSpark SQLfrom pyspark.sql import SparkSession spark SparkSession.builder \ .appName(DailySalesETL) \ .config(spark.sql.warehouse.dir, /data/warehouse) \ .enableHiveSupport() \ .getOrCreate() # 读取 ODS ods spark.sql( SELECT * FROM ods_orders WHERE dt 2024-01-15 ) # 清洗 dwd ods.filter(order_id IS NOT NULL) \ .withColumn(status, upper(col(status))) \ .withColumn(province, ip_to_province(col(ip))) # 写入 DWD dwd.write.mode(overwrite) \ .partitionBy(dt) \ .format(orc) \ .saveAsTable(dwd_orders) # DWS 汇总 spark.sql( INSERT OVERWRITE TABLE dws_user_daily PARTITION (dt2024-01-15) SELECT user_id, COUNT(*) as order_count, SUM(amount) as total_amount FROM dwd_orders WHERE dt 2024-01-15 AND status PAID GROUP BY user_id )Spark Streaming微批# Structured Streaming from pyspark.sql.functions import * # 读取 Kafka stream spark \ .readStream \ .format(kafka) \ .option(kafka.bootstrap.servers, kafka:9092) \ .option(subscribe, orders) \ .load() # 解析 orders stream.select( from_json(col(value).cast(string), schema).alias(data) ).select(data.*) # 聚合 result orders \ .withWatermark(create_time, 5 minutes) \ .groupBy( window(col(create_time), 5 minutes), col(status) ) \ .agg(count(*).alias(order_count)) # 写入 result.writeStream \ .format(console) \ .outputMode(append) \ .start() \ .awaitTermination()Spark 调优# 关键参数 spark.conf.set(spark.sql.shuffle.partitions, 200) # Shuffle 分区数 spark.conf.set(spark.sql.adaptive.enabled, true) # AQE 自适应 spark.conf.set(spark.sql.adaptive.coalescePartitions.enabled, true) spark.conf.set(spark.serializer, org.apache.spark.serializer.KryoSerializer) spark.conf.set(spark.sql.broadcastTimeout, 600) # 广播 JOIN spark.conf.set(spark.sql.autoBroadcastJoinThreshold, 10485760) # 10MB四、Flink 实时计算架构Flink JobManager ├── Dispatcher └── ResourceManager TaskManager多节点 ├── Task Slot 1 → Source ├── Task Slot 2 → Map └── Task Slot 3 → Sink实时数仓// Flink 实时数仓Kafka → Flink → Kafka多层 // 1. ODS消费原始数据 DataStreamString odsStream env.fromSource( KafkaSource.Stringbuilder() .setBootstrapServers(kafka:9092) .setTopics(ods-orders) .setStartingOffsets(OffsetsInitializer.earliest()) .setValueOnlyDeserializer(new SimpleStringSchema()) .build(), WatermarkStrategy.forBoundedOutOfOrderness(Duration.ofMinutes(1)), ods-source ); // 2. DWD清洗 DataStreamJSONObject dwdStream odsStream .map(JSON::parseObject) .filter(obj - obj.containsKey(order_id)) // 过滤脏数据 .map(obj - { obj.put(status, obj.getString(status).toUpperCase()); return obj; }); // 3. DWS实时聚合 DataStreamTuple2String, Double dwsStream dwdStream .keyBy(obj - obj.getString(user_id)) .window(TumblingEventTimeWindows.of(Time.minutes(5))) .aggregate(new AggregateFunctionJSONObject, Tuple2Integer, Double, Tuple2Integer, Double() { Override public Tuple2Integer, Double createAccumulator() { return Tuple2.of(0, 0.0); } Override public Tuple2Integer, Double add(JSONObject value, Tuple2Integer, Double acc) { return Tuple2.of(acc.f0 1, acc.f1 value.getDouble(amount)); } Override public Tuple2Integer, Double getResult(Tuple2Integer, Double acc) { return acc; } Override public Tuple2Integer, Double merge(Tuple2Integer, Double a, Tuple2Integer, Double b) { return Tuple2.of(a.f0 b.f0, a.f1 b.f1); } }); // 4. 写入 Kafka DWS Topic dwsStream .map(value - value.f0 , value.f1) .sinkTo(KafkaSink.Stringbuilder() .setBootstrapServers(kafka:9092) .setRecordSerializer( KafkaRecordSerializationSchema.builder() .setTopic(dws-user-5min) .setValueSerializationSchema(new SimpleStringSchema()) .build() ) .setDeliverGuarantee(DeliveryGuarantee.EXACTLY_ONCE) .build() ); env.execute(Real-time Data Warehouse);Flink SQL-- 创建 Kafka 源表 CREATE TABLE orders_kafka ( order_id BIGINT, user_id BIGINT, amount DECIMAL(10,2), status STRING, create_time TIMESTAMP(3), WATERMARK FOR create_time AS create_time - INTERVAL 1 MINUTE ) WITH ( connector kafka, topic ods-orders, properties.bootstrap.servers kafka:9092, format json, scan.startup.mode latest-offset ); -- 创建 ClickHouse Sink 表 CREATE TABLE orders_clickhouse ( dt DATE, total_orders BIGINT, total_amount DECIMAL(15,2) ) WITH ( connector clickhouse, url clickhouse:8123, database-name realtime, table-name daily_sales ); -- 实时聚合写入 INSERT INTO orders_clickhouse SELECT DATE(create_time) AS dt, COUNT(*) AS total_orders, SUM(amount) AS total_amount FROM orders_kafka WHERE status PAID GROUP BY DATE(create_time), TUMBLE(create_time, INTERVAL 1 DAY);Flink 状态管理// Keyed State public class OrderCounter extends KeyedProcessFunctionString, Order, Result { private ValueStateLong countState; private ValueStateDouble sumState; Override public void open(Configuration parameters) { ValueStateDescriptorLong countDesc new ValueStateDescriptor(count, Long.class); countState getRuntimeContext().getState(countDesc); ValueStateDescriptorDouble sumDesc new ValueStateDescriptor(sum, Double.class); sumState getRuntimeContext().getState(sumDesc); } Override public void processElement(Order order, Context ctx, CollectorResult out) { long count countState.value() ! null ? countState.value() : 0; double sum sumState.value() ! null ? sumState.value() : 0; countState.update(count 1); sumState.update(sum order.getAmount()); out.collect(new Result(order.getUserId(), count 1, sum order.getAmount())); } }五、ClickHouse 实时分析部署# Docker docker run -d \ --name clickhouse \ -p 8123:8123 \ -p 9000:9000 \ -v /data/clickhouse:/var/lib/clickhouse \ clickhouse/clickhouse-server:24.3建表-- MergeTree 引擎OLAP 标准引擎 CREATE TABLE realtime.daily_sales ( dt Date, total_orders UInt64, total_amount Decimal(15,2), paying_users UInt64 ) ENGINE MergeTree() PARTITION BY toYYYYMM(dt) ORDER BY dt; -- ReplacingMergeTree去重 CREATE TABLE realtime.orders ( order_id UInt64, user_id UInt64, amount Decimal(10,2), status String, create_time DateTime ) ENGINE ReplacingMergeTree(create_time) PARTITION BY toDate(create_time) ORDER BY (order_id); -- AggregatingMergeTree预聚合 CREATE TABLE realtime.user_agg ( user_id UInt64, order_count AggregateFunction(count, UInt64), total_amount AggregateFunction(sum, Decimal(10,2)) ) ENGINE AggregatingMergeTree() ORDER BY user_id;实时写入-- 从 Flink 写入 -- 或直接 Kafka → ClickHouse CREATE TABLE realtime.kafka_orders ( order_id UInt64, user_id UInt64, amount Decimal(10,2), status String, create_time DateTime ) ENGINE Kafka() SETTINGS kafka_broker_list kafka:9092, kafka_topic_list ods-orders, kafka_group_name clickhouse-consumer, kafka_format JSONEachRow; -- 物化视图自动消费 Kafka 写入 MergeTree CREATE MATERIALIZED VIEW realtime.orders_mv TO realtime.orders AS SELECT * FROM realtime.kafka_orders;查询优化-- 列裁剪只查需要的列 SELECT user_id, sum(amount) FROM realtime.orders WHERE create_time now() - INTERVAL 1 DAY GROUP BY user_id; -- 分区过滤 SELECT count() FROM realtime.orders WHERE dt 2024-01-15; -- 只扫一个分区 -- 聚合下推 SELECT status, count(), sum(amount) FROM realtime.orders WHERE dt 2024-01-01 GROUP BY status;六、引擎选型总结场景引擎理由离线数仓 ETLHive/Spark SQLSQL 友好、批量处理实时数仓Flink Kafka低延迟、流处理实时 OLAP 查询ClickHouse毫秒级查询大数据批处理Spark通用、灵活全文检索Elasticsearch倒排索引高频点查HBase低延迟实时指标Redis内存级⚠️踩坑提示- Hive/Spark Join 大表用广播 Join- Flink Checkpoint 间隔不要太短1-5 分钟- ClickHouse 避免高频小批量写入用 Buffer 表- Spark AQE 开启后大幅提升性能要点回顾Hive 用于数仓 SQL 查询ODS→DWD→DWS→ADS 分层Spark 通用计算引擎离线批处理首选Flink 实时流处理低延迟、Exactly-OnceClickHouse 实时 OLAP毫秒级查询数仓分层 选取合适引擎是性能的关键下一篇预告下一篇【中台·数据篇】数据治理元数据管理、数据血缘与数据质量监控将讲解数据治理。
返回列表