全部笔记All notes

Flink + Kafka、HBase、Iceberg 实时数据湖构建指南

阅读 8m 00s8m 00s read

概述

在现代大数据架构中,Flink + Kafka + HBase + Iceberg 的组合提供了完整的实时数据处理和存储解决方案:

  • Kafka: 高吞吐量的分布式消息队列,作为数据总线
  • Flink: 流批一体的计算引擎,提供实时处理能力
  • HBase: 高性能的 NoSQL 数据库,支持实时查询
  • Iceberg: 现代化的数据湖表格式,支持 ACID 事务和 Schema 演进

💡 提示: 这种架构结合了流处理的实时性和数据湖的灵活性,适用于需要同时支持实时分析和历史数据查询的场景。

核心优势

组件作用优势
Kafka数据缓冲和传输高吞吐、解耦、持久化
Flink实时计算处理低延迟、精确一次、状态管理
HBase实时数据存储毫秒级查询、随机读写
Iceberg数据湖存储ACID 事务、Schema 演进、时间旅行

架构设计

Lambda 架构

Lambda 架构同时运行批处理和流处理系统:

                          ┌─────────────┐
                          │   Kafka     │
                          │ (数据总线)   │
                          └──────┬──────┘
                                 │
                      ┌──────────┴──────────┐
                      │                     │
                      ▼                     ▼
               ┌──────────┐          ┌──────────┐
               │  Speed   │          │  Batch   │
               │  Layer   │          │  Layer   │
               │ (Flink)  │          │ (Spark)  │
               └─────┬────┘          └─────┬────┘
                     │                     │
                     ▼                     ▼
               ┌──────────┐          ┌──────────┐
               │  HBase   │          │ Iceberg  │
               │(实时视图) │          │(批量视图) │
               └──────────┘          └──────────┘
                     │                     │
                     └──────────┬──────────┘
                                │
                          ┌─────▼─────┐
                          │  Serving  │
                          │   Layer   │
                          └───────────┘

Kappa 架构

Kappa 架构仅使用流处理系统:

┌─────────────┐     ┌─────────────┐     ┌─────────────┐
│   Kafka     │────▶│   Flink     │────▶│  HBase/     │
│ (数据源)     │     │ (流处理引擎) │     │  Iceberg    │
└─────────────┘     └─────────────┘     └─────────────┘
                           │
                           ▼
                    ┌─────────────┐
                    │ Checkpoint  │
                    │  & State    │
                    └─────────────┘

湖仓一体架构

结合数据湖和数据仓库的优势:

public class LakeHouseArchitecture {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        
        // 配置 Checkpoint
        env.enableCheckpointing(60000);
        env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);
        
        // 从 Kafka 读取数据
        KafkaSource<Event> kafkaSource = KafkaSource.<Event>builder()
            .setBootstrapServers("localhost:9092")
            .setTopics("events")
            .setGroupId("flink-consumer")
            .setStartingOffsets(OffsetsInitializer.earliest())
            .setValueOnlyDeserializer(new EventDeserializer())
            .build();
        
        DataStream<Event> events = env.fromSource(
            kafkaSource, 
            WatermarkStrategy.forBoundedOutOfOrderness(Duration.ofSeconds(5)),
            "Kafka Source"
        );
        
        // 实时层:写入 HBase
        events.addSink(new HBaseSink());
        
        // 批处理层:写入 Iceberg
        events.addSink(new IcebergSink());
        
        env.execute("Lake House Architecture");
    }
}

Kafka 集成

Kafka Source

基础配置

import org.apache.flink.connector.kafka.source.KafkaSource;
import org.apache.flink.connector.kafka.source.enumerator.initializer.OffsetsInitializer;

public class KafkaSourceExample {
    public static KafkaSource<String> createKafkaSource() {
        return KafkaSource.<String>builder()
            .setBootstrapServers("broker1:9092,broker2:9092")
            .setTopics("user-events", "system-events")
            .setGroupId("flink-consumer-group")
            .setStartingOffsets(OffsetsInitializer.earliest())
            .setValueOnlyDeserializer(new SimpleStringSchema())
            // 动态分区发现
            .setProperty("partition.discovery.interval.ms", "10000")
            // 消费者配置
            .setProperty("enable.auto.commit", "false")
            .setProperty("max.poll.records", "500")
            .build();
    }
}

自定义反序列化

public class JsonEventDeserializer implements DeserializationSchema<Event> {
    private static final ObjectMapper objectMapper = new ObjectMapper();
    
    @Override
    public Event deserialize(byte[] message) throws IOException {
        return objectMapper.readValue(message, Event.class);
    }
    
    @Override
    public boolean isEndOfStream(Event nextElement) {
        return false;
    }
    
    @Override
    public TypeInformation<Event> getProducedType() {
        return TypeInformation.of(Event.class);
    }
}

带 Schema Registry 的 Avro 反序列化

public class AvroKafkaSource {
    public static KafkaSource<GenericRecord> createAvroSource() {
        ConfluentRegistryAvroDeserializationSchema<GenericRecord> schema = 
            ConfluentRegistryAvroDeserializationSchema.forGeneric(
                userSchema,
                "http://schema-registry:8081"
            );
        
        return KafkaSource.<GenericRecord>builder()
            .setBootstrapServers("localhost:9092")
            .setTopics("avro-events")
            .setGroupId("flink-avro-consumer")
            .setValueOnlyDeserializer(schema)
            .build();
    }
}

Kafka Sink

基础 Sink 配置

public class KafkaSinkExample {
    public static KafkaSink<String> createKafkaSink() {
        return KafkaSink.<String>builder()
            .setBootstrapServers("localhost:9092")
            .setRecordSerializer(KafkaRecordSerializationSchema.builder()
                .setTopic("output-topic")
                .setValueSerializationSchema(new SimpleStringSchema())
                .build()
            )
            .setDeliverGuarantee(DeliveryGuarantee.EXACTLY_ONCE)
            .setTransactionalIdPrefix("flink-producer")
            .build();
    }
}

动态路由 Sink

public class DynamicRoutingKafkaSink {
    public static KafkaSink<Event> createDynamicSink() {
        return KafkaSink.<Event>builder()
            .setBootstrapServers("localhost:9092")
            .setRecordSerializer(new KafkaRecordSerializationSchema<Event>() {
                @Override
                public ProducerRecord<byte[], byte[]> serialize(
                        Event event, KafkaSinkContext context, Long timestamp) {
                    // 根据事件类型动态路由
                    String topic = "events-" + event.getType().toLowerCase();
                    byte[] key = event.getId().getBytes();
                    byte[] value = JSON.toJSONBytes(event);
                    
                    return new ProducerRecord<>(topic, null, timestamp, key, value);
                }
            })
            .setDeliverGuarantee(DeliveryGuarantee.AT_LEAST_ONCE)
            .build();
    }
}

Exactly-Once 语义

端到端 Exactly-Once 配置

public class ExactlyOnceKafkaExample {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        
        // 启用 Checkpoint
        env.enableCheckpointing(5000);
        env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);
        env.getCheckpointConfig().setMinPauseBetweenCheckpoints(500);
        
        // Kafka Source
        KafkaSource<Transaction> source = KafkaSource.<Transaction>builder()
            .setBootstrapServers("localhost:9092")
            .setTopics("transactions")
            .setGroupId("exactly-once-consumer")
            .setStartingOffsets(OffsetsInitializer.committedOffsets())
            .setValueOnlyDeserializer(new TransactionDeserializer())
            .build();
        
        DataStream<Transaction> transactions = env.fromSource(
            source,
            WatermarkStrategy.forMonotonousTimestamps(),
            "Kafka Source"
        );
        
        // 处理逻辑
        DataStream<TransactionResult> results = transactions
            .keyBy(Transaction::getUserId)
            .process(new TransactionProcessor());
        
        // Kafka Sink with Exactly-Once
        KafkaSink<TransactionResult> sink = KafkaSink.<TransactionResult>builder()
            .setBootstrapServers("localhost:9092")
            .setRecordSerializer(KafkaRecordSerializationSchema.builder()
                .setTopic("transaction-results")
                .setValueSerializationSchema(new TransactionResultSerializer())
                .build()
            )
            .setDeliverGuarantee(DeliveryGuarantee.EXACTLY_ONCE)
            .setTransactionalIdPrefix("transaction-processor")
            .setProperty("transaction.timeout.ms", "900000") // 15分钟
            .build();
        
        results.sinkTo(sink);
        
        env.execute("Exactly-Once Kafka Pipeline");
    }
}

HBase 集成

HBase Sink

基础 HBase Sink

public class HBaseSinkExample {
    public static class HBaseSinkFunction extends RichSinkFunction<Event> {
        private transient Connection connection;
        private transient BufferedMutator mutator;
        
        @Override
        public void open(Configuration parameters) throws Exception {
            org.apache.hadoop.conf.Configuration config = HBaseConfiguration.create();
            config.set("hbase.zookeeper.quorum", "localhost");
            config.set("hbase.zookeeper.property.clientPort", "2181");
            
            connection = ConnectionFactory.createConnection(config);
            
            BufferedMutatorParams params = new BufferedMutatorParams(TableName.valueOf("events"))
                .writeBufferSize(4 * 1024 * 1024) // 4MB
                .maxKeyValueSize(10 * 1024 * 1024); // 10MB
                
            mutator = connection.getBufferedMutator(params);
        }
        
        @Override
        public void invoke(Event event, Context context) throws Exception {
            Put put = new Put(Bytes.toBytes(event.getId()));
            put.addColumn(
                Bytes.toBytes("cf"),
                Bytes.toBytes("data"),
                Bytes.toBytes(JSON.toJSONString(event))
            );
            put.addColumn(
                Bytes.toBytes("cf"),
                Bytes.toBytes("timestamp"),
                Bytes.toBytes(event.getTimestamp())
            );
            
            mutator.mutate(put);
        }
        
        @Override
        public void close() throws Exception {
            if (mutator != null) {
                mutator.close();
            }
            if (connection != null) {
                connection.close();
            }
        }
    }
}

异步 HBase Sink

public class AsyncHBaseSink extends RichAsyncFunction<Event, String> {
    private transient AsyncConnection asyncConnection;
    private transient AsyncTable<AdvancedScanResultConsumer> table;
    
    @Override
    public void open(Configuration parameters) throws Exception {
        org.apache.hadoop.conf.Configuration config = HBaseConfiguration.create();
        config.set("hbase.zookeeper.quorum", "localhost");
        
        asyncConnection = ConnectionFactory.createAsyncConnection(config).get();
        table = asyncConnection.getTable(TableName.valueOf("events"));
    }
    
    @Override
    public void asyncInvoke(Event event, ResultFuture<String> resultFuture) {
        Put put = new Put(Bytes.toBytes(event.getId()));
        put.addColumn(Bytes.toBytes("cf"), Bytes.toBytes("data"), 
                      Bytes.toBytes(JSON.toJSONString(event)));
        
        CompletableFuture<Void> future = table.put(put);
        
        future.whenComplete((result, error) -> {
            if (error != null) {
                resultFuture.completeExceptionally(error);
            } else {
                resultFuture.complete(Collections.singletonList("Success: " + event.getId()));
            }
        });
    }
    
    @Override
    public void close() throws Exception {
        if (asyncConnection != null) {
            asyncConnection.close();
        }
    }
}

HBase Lookup Join

维表关联实现

public class HBaseLookupFunction extends TableFunction<Row> {
    private final String tableName;
    private final String family;
    private final String qualifier;
    private Connection connection;
    private Table table;
    
    public HBaseLookupFunction(String tableName, String family, String qualifier) {
        this.tableName = tableName;
        this.family = family;
        this.qualifier = qualifier;
    }
    
    @Override
    public void open(FunctionContext context) throws Exception {
        org.apache.hadoop.conf.Configuration config = HBaseConfiguration.create();
        connection = ConnectionFactory.createConnection(config);
        table = connection.getTable(TableName.valueOf(tableName));
    }
    
    public void eval(String rowKey) {
        try {
            Get get = new Get(Bytes.toBytes(rowKey));
            get.addColumn(Bytes.toBytes(family), Bytes.toBytes(qualifier));
            
            Result result = table.get(get);
            if (!result.isEmpty()) {
                byte[] value = result.getValue(Bytes.toBytes(family), Bytes.toBytes(qualifier));
                collect(Row.of(rowKey, Bytes.toString(value)));
            }
        } catch (IOException e) {
            // 错误处理
        }
    }
    
    @Override
    public void close() throws Exception {
        if (table != null) table.close();
        if (connection != null) connection.close();
    }
}

批量操作优化

public class BatchHBaseSink extends RichSinkFunction<List<Event>> 
        implements CheckpointedFunction {
    
    private transient ListState<Event> checkpointedState;
    private final List<Event> bufferedEvents = new ArrayList<>();
    private final int batchSize = 1000;
    private Connection connection;
    private BufferedMutator mutator;
    
    @Override
    public void invoke(List<Event> events, Context context) throws Exception {
        bufferedEvents.addAll(events);
        
        if (bufferedEvents.size() >= batchSize) {
            flush();
        }
    }
    
    private void flush() throws Exception {
        List<Put> puts = new ArrayList<>();
        
        for (Event event : bufferedEvents) {
            Put put = new Put(Bytes.toBytes(event.getId()));
            put.addColumn(Bytes.toBytes("cf"), Bytes.toBytes("data"), 
                         Bytes.toBytes(JSON.toJSONString(event)));
            puts.add(put);
        }
        
        mutator.mutate(puts);
        mutator.flush();
        bufferedEvents.clear();
    }
    
    @Override
    public void snapshotState(FunctionSnapshotContext context) throws Exception {
        checkpointedState.clear();
        for (Event event : bufferedEvents) {
            checkpointedState.add(event);
        }
    }
    
    @Override
    public void initializeState(FunctionInitializationContext context) throws Exception {
        ListStateDescriptor<Event> descriptor = new ListStateDescriptor<>(
            "buffered-events",
            TypeInformation.of(Event.class)
        );
        
        checkpointedState = context.getOperatorStateStore().getListState(descriptor);
        
        if (context.isRestored()) {
            for (Event event : checkpointedState.get()) {
                bufferedEvents.add(event);
            }
        }
    }
}

Iceberg 集成

流式写入

DataStream API 写入 Iceberg

public class IcebergStreamingSink {
    public static void writeToIceberg(DataStream<RowData> stream) {
        // 创建 Iceberg 表
        TableIdentifier tableId = TableIdentifier.of("db", "events");
        
        // 配置 Iceberg Sink
        FlinkSink.forRowData(stream)
            .table(table)
            .tableLoader(TableLoader.fromHadoopTable("hdfs://namenode:9000/warehouse/db/events"))
            .writeParallelism(4)
            .build();
    }
    
    // 使用 SQL API
    public static void writeWithSQL(StreamTableEnvironment tableEnv) {
        tableEnv.executeSql(
            "CREATE CATALOG iceberg_catalog WITH (" +
            "  'type' = 'iceberg'," +
            "  'catalog-type' = 'hive'," +
            "  'uri' = 'thrift://localhost:9083'," +
            "  'warehouse' = 'hdfs://namenode:9000/warehouse'" +
            ")"
        );
        
        tableEnv.executeSql(
            "CREATE TABLE IF NOT EXISTS iceberg_catalog.db.events (" +
            "  id STRING," +
            "  user_id STRING," +
            "  event_type STRING," +
            "  event_time TIMESTAMP(3)," +
            "  properties MAP<STRING, STRING>," +
            "  PRIMARY KEY (id) NOT ENFORCED" +
            ") PARTITIONED BY (event_type, days(event_time)) WITH (" +
            "  'format-version' = '2'," +
            "  'write.format.default' = 'parquet'," +
            "  'write.parquet.compression-codec' = 'snappy'" +
            ")"
        );
        
        tableEnv.executeSql(
            "INSERT INTO iceberg_catalog.db.events " +
            "SELECT * FROM kafka_events"
        );
    }
}

Iceberg 表配置选项

public class IcebergTableConfiguration {
    public static Table createOptimizedTable() {
        Schema schema = new Schema(
            Types.NestedField.required(1, "id", Types.StringType.get()),
            Types.NestedField.required(2, "user_id", Types.StringType.get()),
            Types.NestedField.required(3, "event_type", Types.StringType.get()),
            Types.NestedField.required(4, "event_time", Types.TimestampType.withZone()),
            Types.NestedField.optional(5, "properties", Types.MapType.ofRequired(6, 7,
                Types.StringType.get(), Types.StringType.get()))
        );
        
        PartitionSpec spec = PartitionSpec.builderFor(schema)
            .identity("event_type")
            .day("event_time")
            .build();
        
        Map<String, String> properties = ImmutableMap.of(
            "write.format.default", "parquet",
            "write.parquet.compression-codec", "snappy",
            "write.metadata.compression-codec", "gzip",
            "write.target-file-size-bytes", String.valueOf(128 * 1024 * 1024), // 128MB
            "write.distribution-mode", "hash"
        );
        
        return catalog.createTable(tableId, schema, spec, properties);
    }
}

批流一体

流批混合处理

public class StreamBatchHybrid {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);
        
        // 创建 Iceberg Catalog
        tableEnv.executeSql(
            "CREATE CATALOG iceberg WITH (" +
            "  'type' = 'iceberg'," +
            "  'catalog-type' = 'hadoop'," +
            "  'warehouse' = 'hdfs://namenode:9000/warehouse'" +
            ")"
        );
        
        // 实时流处理
        tableEnv.executeSql(
            "CREATE TABLE kafka_source (" +
            "  id STRING," +
            "  event_time TIMESTAMP(3)," +
            "  WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND" +
            ") WITH (" +
            "  'connector' = 'kafka'," +
            "  'topic' = 'events'," +
            "  'properties.bootstrap.servers' = 'localhost:9092'," +
            "  'format' = 'json'" +
            ")"
        );
        
        // 流式写入 Iceberg
        tableEnv.executeSql(
            "INSERT INTO iceberg.db.events " +
            "SELECT * FROM kafka_source"
        );
        
        // 批量查询 Iceberg
        Table batchResult = tableEnv.sqlQuery(
            "SELECT " +
            "  DATE_FORMAT(event_time, 'yyyy-MM-dd') as event_date," +
            "  COUNT(*) as event_count " +
            "FROM iceberg.db.events " +
            "WHERE event_time >= CURRENT_TIMESTAMP - INTERVAL '7' DAY " +
            "GROUP BY DATE_FORMAT(event_time, 'yyyy-MM-dd')"
        );
        
        // 输出结果
        tableEnv.toRetractStream(batchResult, Row.class).print();
        
        env.execute("Stream-Batch Hybrid Processing");
    }
}

表维护

Iceberg 表维护操作

public class IcebergTableMaintenance {
    private final Table table;
    private final SparkSession spark;
    
    public void compactFiles() {
        // 小文件合并
        Actions.forTable(spark, table)
            .rewriteDataFiles()
            .filter(Expressions.equal("event_type", "click"))
            .targetSizeInBytes(128 * 1024 * 1024) // 128MB
            .execute();
    }
    
    public void expireSnapshots() {
        // 过期快照清理
        table.expireSnapshots()
            .expireOlderThan(System.currentTimeMillis() - TimeUnit.DAYS.toMillis(7))
            .retainLast(3)
            .commit();
    }
    
    public void removeOrphanFiles() {
        // 孤儿文件清理
        Actions.forTable(spark, table)
            .removeOrphanFiles()
            .olderThan(System.currentTimeMillis() - TimeUnit.DAYS.toMillis(3))
            .execute();
    }
    
    public void rewriteManifests() {
        // 重写 manifest 文件
        table.rewriteManifests()
            .clusterBy(DataFile::partition)
            .commit();
    }
    
    // 定期维护任务
    public static class MaintenanceJob {
        public static void scheduleMaintenance() {
            ScheduledExecutorService scheduler = Executors.newScheduledThreadPool(1);
            
            scheduler.scheduleAtFixedRate(() -> {
                try {
                    IcebergTableMaintenance maintenance = new IcebergTableMaintenance();
                    maintenance.compactFiles();
                    maintenance.expireSnapshots();
                    maintenance.removeOrphanFiles();
                } catch (Exception e) {
                    LOG.error("Maintenance failed", e);
                }
            }, 0, 24, TimeUnit.HOURS);
        }
    }
}

实战案例

实时数据湖架构

public class RealtimeDataLake {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);
        
        // 配置环境
        env.setParallelism(4);
        env.enableCheckpointing(60000);
        env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);
        
        // 1. 定义 Kafka 数据源
        tableEnv.executeSql(
            "CREATE TABLE raw_events (" +
            "  event_id STRING," +
            "  user_id STRING," +
            "  event_type STRING," +
            "  event_time TIMESTAMP(3)," +
            "  attributes MAP<STRING, STRING>," +
            "  WATERMARK FOR event_time AS event_time - INTERVAL '10' SECOND" +
            ") WITH (" +
            "  'connector' = 'kafka'," +
            "  'topic' = 'raw-events'," +
            "  'properties.bootstrap.servers' = 'localhost:9092'," +
            "  'properties.group.id' = 'data-lake-consumer'," +
            "  'format' = 'json'," +
            "  'json.fail-on-missing-field' = 'false'" +
            ")"
        );
        
        // 2. 数据清洗和标准化
        Table cleanedEvents = tableEnv.sqlQuery(
            "SELECT " +
            "  event_id," +
            "  user_id," +
            "  UPPER(event_type) as event_type," +
            "  event_time," +
            "  attributes['device_type'] as device_type," +
            "  attributes['location'] as location," +
            "  CASE " +
            "    WHEN event_type IN ('click', 'view') THEN 'engagement'" +
            "    WHEN event_type IN ('purchase', 'add_to_cart') THEN 'transaction'" +
            "    ELSE 'other' " +
            "  END as event_category " +
            "FROM raw_events " +
            "WHERE user_id IS NOT NULL AND event_id IS NOT NULL"
        );
        
        // 3. 实时聚合
        Table realtimeMetrics = tableEnv.sqlQuery(
            "SELECT " +
            "  window_start," +
            "  window_end," +
            "  event_category," +
            "  COUNT(DISTINCT user_id) as unique_users," +
            "  COUNT(*) as event_count," +
            "  COUNT(DISTINCT device_type) as device_types " +
            "FROM TABLE(" +
            "  TUMBLE(TABLE cleaned_events, DESCRIPTOR(event_time), INTERVAL '5' MINUTE)" +
            ") " +
            "GROUP BY window_start, window_end, event_category"
        );
        
        // 4. 写入 HBase(实时查询层)
        tableEnv.executeSql(
            "CREATE TABLE hbase_realtime_metrics (" +
            "  rowkey STRING," +
            "  cf ROW<" +
            "    unique_users BIGINT," +
            "    event_count BIGINT," +
            "    device_types BIGINT" +
            "  >," +
            "  PRIMARY KEY (rowkey) NOT ENFORCED" +
            ") WITH (" +
            "  'connector' = 'hbase-2.2'," +
            "  'table-name' = 'realtime_metrics'," +
            "  'zookeeper.quorum' = 'localhost:2181'" +
            ")"
        );
        
        tableEnv.executeSql(
            "INSERT INTO hbase_realtime_metrics " +
            "SELECT " +
            "  CONCAT(DATE_FORMAT(window_start, 'yyyyMMddHHmm'), '_', event_category) as rowkey," +
            "  ROW(unique_users, event_count, device_types) as cf " +
            "FROM realtime_metrics"
        );
        
        // 5. 写入 Iceberg(数据湖存储)
        tableEnv.executeSql(
            "INSERT INTO iceberg.data_lake.events " +
            "SELECT * FROM cleaned_events"
        );
        
        // 6. CDC 数据同步(维度表)
        tableEnv.executeSql(
            "CREATE TABLE mysql_users (" +
            "  user_id STRING," +
            "  user_name STRING," +
            "  registration_date DATE," +
            "  user_level INT," +
            "  PRIMARY KEY (user_id) NOT ENFORCED" +
            ") WITH (" +
            "  'connector' = 'mysql-cdc'," +
            "  'hostname' = 'localhost'," +
            "  'port' = '3306'," +
            "  'username' = 'root'," +
            "  'password' = 'password'," +
            "  'database-name' = 'users'," +
            "  'table-name' = 'user_info'" +
            ")"
        );
        
        // 维度关联并写入宽表
        tableEnv.executeSql(
            "INSERT INTO iceberg.data_lake.enriched_events " +
            "SELECT " +
            "  e.*," +
            "  u.user_name," +
            "  u.user_level," +
            "  CURRENT_TIMESTAMP as process_time " +
            "FROM cleaned_events e " +
            "LEFT JOIN mysql_users FOR SYSTEM_TIME AS OF e.event_time AS u " +
            "ON e.user_id = u.user_id"
        );
        
        env.execute("Realtime Data Lake");
    }
}

用户行为分析系统

public class UserBehaviorAnalysis {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        
        // Kafka 数据源
        KafkaSource<UserBehavior> source = KafkaSource.<UserBehavior>builder()
            .setBootstrapServers("localhost:9092")
            .setTopics("user-behavior")
            .setGroupId("behavior-analysis")
            .setStartingOffsets(OffsetsInitializer.latest())
            .setValueOnlyDeserializer(new JsonDeserializationSchema<>(UserBehavior.class))
            .build();
        
        DataStream<UserBehavior> behaviors = env.fromSource(
            source,
            WatermarkStrategy.<UserBehavior>forBoundedOutOfOrderness(Duration.ofSeconds(5))
                .withTimestampAssigner((behavior, timestamp) -> behavior.getTimestamp()),
            "User Behavior Source"
        );
        
        // 1. 实时用户活跃度分析
        DataStream<UserActivity> userActivity = behaviors
            .keyBy(UserBehavior::getUserId)
            .window(SlidingEventTimeWindows.of(Time.minutes(30), Time.minutes(5)))
            .aggregate(new UserActivityAggregator());
        
        // 写入 HBase
        userActivity.addSink(new UserActivityHBaseSink());
        
        // 2. 用户路径分析
        DataStream<UserPath> userPaths = behaviors
            .keyBy(UserBehavior::getSessionId)
            .window(SessionWindows.withGap(Time.minutes(30)))
            .process(new UserPathAnalyzer());
        
        // 3. 实时推荐特征计算
        DataStream<UserFeature> userFeatures = behaviors
            .keyBy(UserBehavior::getUserId)
            .process(new UserFeatureExtractor())
            .uid("user-feature-extractor");
        
        // 异步 IO 查询用户画像
        DataStream<EnrichedUserFeature> enrichedFeatures = AsyncDataStream
            .unorderedWait(
                userFeatures,
                new UserProfileAsyncFunction(),
                60, TimeUnit.SECONDS,
                100
            );
        
        // 4. 批流结合:历史数据回填
        if (args.length > 0 && "backfill".equals(args[0])) {
            DataStream<UserBehavior> historicalData = env
                .readFile(
                    new TextInputFormat(new Path("hdfs://namenode:9000/history")),
                    "hdfs://namenode:9000/history"
                )
                .map(new HistoricalDataParser());
            
            behaviors = behaviors.union(historicalData);
        }
        
        // 5. 写入 Iceberg 数据湖
        FlinkSink.forRowData(
            behaviors.map(new BehaviorToRowDataMapper()),
            TableLoader.fromHadoopTable("hdfs://namenode:9000/warehouse/behavior_lake")
        )
        .tableSchema(TableSchema.builder()
            .field("user_id", DataTypes.STRING())
            .field("behavior", DataTypes.STRING())
            .field("timestamp", DataTypes.TIMESTAMP(3))
            .build())
        .build();
        
        env.execute("User Behavior Analysis");
    }
}

// 用户活跃度聚合器
class UserActivityAggregator implements AggregateFunction<UserBehavior, UserActivityAccumulator, UserActivity> {
    @Override
    public UserActivityAccumulator createAccumulator() {
        return new UserActivityAccumulator();
    }
    
    @Override
    public UserActivityAccumulator add(UserBehavior behavior, UserActivityAccumulator acc) {
        acc.addBehavior(behavior);
        return acc;
    }
    
    @Override
    public UserActivity getResult(UserActivityAccumulator acc) {
        return acc.toUserActivity();
    }
    
    @Override
    public UserActivityAccumulator merge(UserActivityAccumulator a, UserActivityAccumulator b) {
        return a.merge(b);
    }
}

IoT 数据处理平台

public class IoTDataPlatform {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        
        // 设置时间语义
        env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime);
        env.enableCheckpointing(30000);
        
        // 1. 多源数据接入
        DataStream<SensorData> mqttStream = env.addSource(new MqttSource("tcp://mqtt-broker:1883"));
        DataStream<SensorData> kafkaStream = env.addSource(createKafkaSource());
        
        DataStream<SensorData> allSensorData = mqttStream.union(kafkaStream)
            .assignTimestampsAndWatermarks(
                WatermarkStrategy.<SensorData>forBoundedOutOfOrderness(Duration.ofSeconds(10))
                    .withTimestampAssigner((data, timestamp) -> data.getTimestamp())
            );
        
        // 2. 数据质量检查
        OutputTag<SensorData> invalidDataTag = new OutputTag<SensorData>("invalid-data"){};
        
        SingleOutputStreamOperator<SensorData> validatedData = allSensorData
            .process(new ProcessFunction<SensorData, SensorData>() {
                @Override
                public void processElement(SensorData data, Context ctx, Collector<SensorData> out) {
                    if (isValid(data)) {
                        out.collect(data);
                    } else {
                        ctx.output(invalidDataTag, data);
                    }
                }
                
                private boolean isValid(SensorData data) {
                    return data.getValue() >= -50 && data.getValue() <= 150 &&
                           data.getDeviceId() != null && !data.getDeviceId().isEmpty();
                }
            });
        
        // 异常数据处理
        DataStream<SensorData> invalidData = validatedData.getSideOutput(invalidDataTag);
        invalidData.addSink(new ErrorDataSink());
        
        // 3. 实时聚合计算
        DataStream<DeviceMetrics> deviceMetrics = validatedData
            .keyBy(SensorData::getDeviceId)
            .window(TumblingEventTimeWindows.of(Time.minutes(1)))
            .aggregate(new DeviceMetricsAggregator());
        
        // 4. 异常检测
        DataStream<Alert> alerts = deviceMetrics
            .keyBy(DeviceMetrics::getDeviceId)
            .process(new AnomalyDetector())
            .filter(alert -> alert.getSeverity() >= Alert.Severity.WARNING);
        
        // 5. 写入实时存储(HBase)
        deviceMetrics.addSink(new HBaseDeviceMetricsSink());
        alerts.addSink(new HBaseAlertSink());
        
        // 6. 写入数据湖(Iceberg)
        DataStream<Row> rowStream = validatedData
            .map(new SensorDataToRowMapper())
            .returns(Types.ROW(
                Types.STRING,  // device_id
                Types.DOUBLE,  // value
                Types.LONG,    // timestamp
                Types.STRING   // location
            ));
        
        Table sensorTable = tableEnv.fromDataStream(rowStream, 
            $("device_id"), $("value"), $("timestamp"), $("location"));
        
        tableEnv.executeSql(
            "INSERT INTO iceberg.iot.sensor_data " +
            "PARTITION (day) " +
            "SELECT *, DATE_FORMAT(FROM_UNIXTIME(timestamp/1000), 'yyyy-MM-dd') as day " +
            "FROM sensor_data"
        );
        
        // 7. 规则引擎集成
        BroadcastStream<Rule> ruleStream = env
            .addSource(new RuleSource())
            .broadcast(new MapStateDescriptor<>("rules", String.class, Rule.class));
        
        DataStream<Alert> ruleAlerts = validatedData
            .connect(ruleStream)
            .process(new DynamicRuleProcessor());
        
        env.execute("IoT Data Platform");
    }
}

// 设备指标聚合器
class DeviceMetricsAggregator implements AggregateFunction<SensorData, MetricsAccumulator, DeviceMetrics> {
    @Override
    public MetricsAccumulator createAccumulator() {
        return new MetricsAccumulator();
    }
    
    @Override
    public MetricsAccumulator add(SensorData data, MetricsAccumulator acc) {
        acc.add(data.getValue());
        return acc;
    }
    
    @Override
    public DeviceMetrics getResult(MetricsAccumulator acc) {
        return new DeviceMetrics(
            acc.getDeviceId(),
            acc.getAvg(),
            acc.getMin(),
            acc.getMax(),
            acc.getCount(),
            acc.getWindowStart(),
            acc.getWindowEnd()
        );
    }
}

// 异常检测器
class AnomalyDetector extends KeyedProcessFunction<String, DeviceMetrics, Alert> {
    private ValueState<DeviceBaseline> baselineState;
    
    @Override
    public void open(Configuration parameters) {
        ValueStateDescriptor<DeviceBaseline> descriptor = new ValueStateDescriptor<>(
            "baseline",
            TypeInformation.of(DeviceBaseline.class)
        );
        baselineState = getRuntimeContext().getState(descriptor);
    }
    
    @Override
    public void processElement(DeviceMetrics metrics, Context ctx, Collector<Alert> out) throws Exception {
        DeviceBaseline baseline = baselineState.value();
        
        if (baseline == null) {
            baseline = new DeviceBaseline(metrics.getDeviceId());
            baselineState.update(baseline);
        }
        
        // 更新基线
        baseline.update(metrics);
        
        // 检测异常
        if (baseline.isAnomaly(metrics)) {
            Alert alert = new Alert(
                metrics.getDeviceId(),
                "Anomaly detected: " + metrics,
                Alert.Severity.HIGH,
                System.currentTimeMillis()
            );
            out.collect(alert);
        }
        
        baselineState.update(baseline);
    }
}

性能优化

Kafka 优化

public class KafkaOptimization {
    // 生产者优化
    public static Properties getOptimizedProducerProps() {
        Properties props = new Properties();
        props.put("bootstrap.servers", "localhost:9092");
        
        // 批处理优化
        props.put("batch.size", 32768); // 32KB
        props.put("linger.ms", 10); // 10ms延迟
        props.put("compression.type", "lz4");
        
        // 缓冲区优化
        props.put("buffer.memory", 67108864); // 64MB
        props.put("max.in.flight.requests.per.connection", 5);
        
        // 重试策略
        props.put("retries", 3);
        props.put("retry.backoff.ms", 100);
        
        return props;
    }
    
    // 消费者优化
    public static Properties getOptimizedConsumerProps() {
        Properties props = new Properties();
        props.put("bootstrap.servers", "localhost:9092");
        
        // 拉取优化
        props.put("fetch.min.bytes", 1048576); // 1MB
        props.put("fetch.max.wait.ms", 500);
        props.put("max.partition.fetch.bytes", 10485760); // 10MB
        
        // 并发优化
        props.put("max.poll.records", 5000);
        props.put("max.poll.interval.ms", 300000); // 5分钟
        
        // 会话管理
        props.put("session.timeout.ms", 30000);
        props.put("heartbeat.interval.ms", 3000);
        
        return props;
    }
}

HBase 优化

public class HBaseOptimization {
    public static Configuration getOptimizedConfig() {
        Configuration config = HBaseConfiguration.create();
        
        // 连接池优化
        config.set("hbase.client.ipc.pool.type", "RoundRobinPool");
        config.set("hbase.client.ipc.pool.size", "10");
        
        // 缓存优化
        config.set("hbase.client.scanner.caching", "1000");
        config.set("hbase.client.scanner.max.result.size", "10485760"); // 10MB
        
        // 重试策略
        config.set("hbase.client.retries.number", "3");
        config.set("hbase.client.pause", "100");
        
        // RPC 超时
        config.set("hbase.rpc.timeout", "60000");
        config.set("hbase.client.operation.timeout", "120000");
        
        return config;
    }
    
    // 批量操作优化
    public static class OptimizedHBaseSink extends RichSinkFunction<Event> {
        private final int batchSize = 1000;
        private final long flushInterval = 5000; // 5秒
        private List<Put> buffer = new ArrayList<>();
        private long lastFlushTime = System.currentTimeMillis();
        
        @Override
        public void invoke(Event event, Context context) throws Exception {
            Put put = createPut(event);
            buffer.add(put);
            
            if (buffer.size() >= batchSize || 
                System.currentTimeMillis() - lastFlushTime > flushInterval) {
                flush();
            }
        }
        
        private void flush() throws IOException {
            if (!buffer.isEmpty()) {
                table.put(buffer);
                buffer.clear();
                lastFlushTime = System.currentTimeMillis();
            }
        }
    }
}

Iceberg 优化

public class IcebergOptimization {
    // 写入优化
    public static Map<String, String> getOptimizedWriteProperties() {
        return ImmutableMap.<String, String>builder()
            // 文件大小优化
            .put("write.target-file-size-bytes", "134217728") // 128MB
            .put("write.parquet.row-group-size-bytes", "134217728")
            .put("write.parquet.page-size-bytes", "1048576") // 1MB
            
            // 并发优化
            .put("write.distribution-mode", "hash")
            .put("write.wap.enabled", "true")
            
            // 压缩优化
            .put("write.parquet.compression-codec", "snappy")
            .put("write.metadata.compression-codec", "gzip")
            
            // 索引优化
            .put("write.parquet.bloom-filter-enabled.column.user_id", "true")
            .put("write.parquet.bloom-filter-enabled.column.event_type", "true")
            
            .build();
    }
    
    // 读取优化
    public static void optimizeTableScan(Table table) {
        table.newScan()
            .caseSensitive(false)
            .select("user_id", "event_type", "event_time")
            .filter(Expressions.greaterThanOrEqual("event_time", 
                    System.currentTimeMillis() - TimeUnit.DAYS.toMillis(7)))
            .planWith(Executors.newFixedThreadPool(4))
            .option("split-size", "134217728") // 128MB
            .option("split-lookback", "10")
            .option("split-open-file-cost", "4194304"); // 4MB
    }
}
public class FlinkJobOptimization {
    public static void optimizeJob(StreamExecutionEnvironment env) {
        // 并行度优化
        env.setParallelism(env.getMaxParallelism());
        
        // 内存优化
        Configuration config = new Configuration();
        config.set(TaskManagerOptions.MANAGED_MEMORY_SIZE, MemorySize.parse("2g"));
        config.set(TaskManagerOptions.NETWORK_MEMORY_FRACTION, 0.15);
        config.set(TaskManagerOptions.TASK_HEAP_MEMORY, MemorySize.parse("1g"));
        
        // 状态后端优化
        env.setStateBackend(new RocksDBStateBackend("hdfs://namenode:9000/flink/checkpoints"));
        
        // Checkpoint 优化
        env.enableCheckpointing(60000);
        env.getCheckpointConfig().setCheckpointTimeout(120000);
        env.getCheckpointConfig().setMaxConcurrentCheckpoints(1);
        env.getCheckpointConfig().enableUnalignedCheckpoints();
        
        // 网络优化
        env.getConfig().setAutoWatermarkInterval(200);
        env.setBufferTimeout(100);
    }
}

监控与运维

指标监控

public class MetricsMonitoring {
    // 自定义指标
    public static class MonitoredKafkaSource extends RichSourceFunction<Event> {
        private transient Counter eventsCounter;
        private transient Meter eventsMeter;
        private transient Histogram lagHistogram;
        
        @Override
        public void open(Configuration parameters) {
            MetricGroup metricGroup = getRuntimeContext().getMetricGroup();
            
            eventsCounter = metricGroup.counter("events_consumed");
            eventsMeter = metricGroup.meter("events_rate", new MeterView(60));
            lagHistogram = metricGroup.histogram("consumer_lag", 
                new DescriptiveStatisticsHistogram(1000));
        }
        
        @Override
        public void run(SourceContext<Event> ctx) throws Exception {
            while (isRunning) {
                ConsumerRecords<String, Event> records = consumer.poll(Duration.ofMillis(100));
                
                for (ConsumerRecord<String, Event> record : records) {
                    Event event = record.value();
                    ctx.collect(event);
                    
                    // 更新指标
                    eventsCounter.inc();
                    eventsMeter.markEvent();
                    
                    long lag = System.currentTimeMillis() - record.timestamp();
                    lagHistogram.update(lag);
                }
            }
        }
    }
    
    // Prometheus 集成
    public static void configurePrometheusReporter(Configuration config) {
        config.setString("metrics.reporters", "prometheus");
        config.setString("metrics.reporter.prometheus.class", 
            "org.apache.flink.metrics.prometheus.PrometheusReporter");
        config.setString("metrics.reporter.prometheus.port", "9249");
    }
}

健康检查

public class HealthCheck {
    public static class SystemHealthMonitor extends ProcessFunction<Event, HealthStatus> {
        private transient ValueState<ComponentHealth> kafkaHealth;
        private transient ValueState<ComponentHealth> hbaseHealth;
        private transient ValueState<ComponentHealth> icebergHealth;
        
        @Override
        public void processElement(Event event, Context ctx, Collector<HealthStatus> out) {
            // 定期健康检查
            ctx.timerService().registerProcessingTimeTimer(
                ctx.timerService().currentProcessingTime() + 60000); // 1分钟
        }
        
        @Override
        public void onTimer(long timestamp, OnTimerContext ctx, Collector<HealthStatus> out) {
            HealthStatus status = new HealthStatus();
            
            // 检查 Kafka
            status.setKafkaHealth(checkKafkaHealth());
            
            // 检查 HBase
            status.setHbaseHealth(checkHBaseHealth());
            
            // 检查 Iceberg
            status.setIcebergHealth(checkIcebergHealth());
            
            out.collect(status);
            
            // 告警
            if (!status.isHealthy()) {
                sendAlert(status);
            }
        }
    }
}

最佳实践

1. 数据分层设计

// ODS层:原始数据层
tableEnv.executeSql("CREATE DATABASE IF NOT EXISTS ods");

// DWD层:明细数据层
tableEnv.executeSql("CREATE DATABASE IF NOT EXISTS dwd");

// DWS层:汇总数据层
tableEnv.executeSql("CREATE DATABASE IF NOT EXISTS dws");

// ADS层:应用数据层
tableEnv.executeSql("CREATE DATABASE IF NOT EXISTS ads");

2. 错误处理策略

public class ErrorHandlingStrategy {
    // 使用侧输出流处理错误
    OutputTag<ErrorRecord> errorTag = new OutputTag<ErrorRecord>("errors"){};
    
    SingleOutputStreamOperator<ProcessedRecord> mainStream = inputStream
        .process(new ProcessFunction<RawRecord, ProcessedRecord>() {
            @Override
            public void processElement(RawRecord record, Context ctx, 
                                     Collector<ProcessedRecord> out) {
                try {
                    ProcessedRecord processed = process(record);
                    out.collect(processed);
                } catch (Exception e) {
                    ctx.output(errorTag, new ErrorRecord(record, e));
                }
            }
        });
    
    // 错误数据处理
    DataStream<ErrorRecord> errorStream = mainStream.getSideOutput(errorTag);
    errorStream.addSink(new ErrorSink());
}

3. 资源隔离

// 为不同优先级的任务设置不同的资源
env.execute("high-priority-job");

// 使用 Slot 共享组
stream.filter(...).slotSharingGroup("source");
stream.keyBy(...).window(...).slotSharingGroup("compute");
stream.addSink(...).slotSharingGroup("sink");

4. 数据质量保证

public class DataQualityAssurance {
    public static void addQualityChecks(DataStream<Event> stream) {
        // 完整性检查
        stream.filter(event -> event.getUserId() != null && event.getEventId() != null);
        
        // 准确性检查
        stream.filter(event -> event.getTimestamp() > 0 && 
                              event.getTimestamp() < System.currentTimeMillis());
        
        // 一致性检查
        stream.keyBy(Event::getEventId)
              .process(new DuplicateDetector());
        
        // 及时性检查
        stream.assignTimestampsAndWatermarks(
            WatermarkStrategy.<Event>forBoundedOutOfOrderness(Duration.ofMinutes(5))
                .withTimestampAssigner((event, timestamp) -> event.getTimestamp())
        );
    }
}

常见问题

Q1: 如何处理数据延迟?

解答:

  1. 调整 Watermark 策略:增加乱序容忍度
  2. 使用 AllowedLateness:允许迟到数据
  3. 侧输出流处理:收集迟到数据单独处理
stream.keyBy(...)
      .window(...)
      .allowedLateness(Time.minutes(5))
      .sideOutputLateData(lateDataTag)
      .aggregate(...);

Q2: 如何保证 Exactly-Once?

解答:

  1. 启用 Checkpoint
  2. 使用支持事务的 Sink
  3. 配置合适的超时时间
env.enableCheckpointing(5000);
env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);

// Kafka Sink
sink.setDeliverGuarantee(DeliveryGuarantee.EXACTLY_ONCE);

Q3: 如何优化小文件问题?

解答:

  1. Iceberg:定期执行文件合并
  2. HBase:使用批量写入
  3. 调整并行度:减少写入并行度

Q4: 如何处理反压?

解答:

  1. 增加并行度
  2. 优化算子逻辑
  3. 调整缓冲区大小
  4. 使用异步 I/O

Q5: 如何监控作业状态?

解答:

  1. 使用 Flink Web UI
  2. 集成 Prometheus + Grafana
  3. 自定义指标和告警

相关文章

高级特性

集成技术

实战案例


总结

Flink + Kafka + HBase + Iceberg 的组合提供了完整的实时数据处理解决方案:

🎯 架构优势

  1. Kafka 提供可靠的数据传输和缓冲
  2. Flink 提供强大的流处理能力
  3. HBase 支持低延迟的实时查询
  4. Iceberg 实现高效的数据湖存储

✅ 最佳实践

  • 合理设计数据分层架构
  • 实现端到端的 Exactly-Once 语义
  • 定期进行表维护和优化
  • 建立完善的监控和告警机制
  • 做好容错和故障恢复准备

🚀 适用场景

  • 实时数据仓库
  • 用户行为分析
  • IoT 数据处理
  • 实时推荐系统
  • 金融风控系统

🚫 注意事项:

  • 合理控制并行度,避免资源浪费
  • 注意数据倾斜问题
  • 定期清理过期数据
  • 监控各组件的健康状态

通过合理的架构设计和优化,这套技术栈可以支撑 PB 级数据的实时处理需求,为企业的数据驱动决策提供强有力的支持。

上一章 / 下一章