概述
在现代大数据架构中,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
}
}
Flink 作业优化
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: 如何处理数据延迟?
解答:
- 调整 Watermark 策略:增加乱序容忍度
- 使用 AllowedLateness:允许迟到数据
- 侧输出流处理:收集迟到数据单独处理
stream.keyBy(...)
.window(...)
.allowedLateness(Time.minutes(5))
.sideOutputLateData(lateDataTag)
.aggregate(...);
Q2: 如何保证 Exactly-Once?
解答:
- 启用 Checkpoint
- 使用支持事务的 Sink
- 配置合适的超时时间
env.enableCheckpointing(5000);
env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);
// Kafka Sink
sink.setDeliverGuarantee(DeliveryGuarantee.EXACTLY_ONCE);
Q3: 如何优化小文件问题?
解答:
- Iceberg:定期执行文件合并
- HBase:使用批量写入
- 调整并行度:减少写入并行度
Q4: 如何处理反压?
解答:
- 增加并行度
- 优化算子逻辑
- 调整缓冲区大小
- 使用异步 I/O
Q5: 如何监控作业状态?
解答:
- 使用 Flink Web UI
- 集成 Prometheus + Grafana
- 自定义指标和告警
相关文章
Flink 核心
- Flink 架构与工作流程 - Flink 核心原理
- Flink DataStream API - DataStream 基础
- Flink DataStream API 高级用法 - 高级特性
高级特性
- Checkpoint & Savepoint - 容错机制
- Flink Table & SQL API - SQL 开发
- Flink CDC - 数据同步
- Flink 性能优化 - 调优指南
- Flink 监控 - 监控方案
集成技术
- Kafka 简明教程 - Kafka 详解
- HBase 教程 - HBase 使用
- Apache Iceberg - Iceberg 官方文档
实战案例
总结
Flink + Kafka + HBase + Iceberg 的组合提供了完整的实时数据处理解决方案:
🎯 架构优势
- Kafka 提供可靠的数据传输和缓冲
- Flink 提供强大的流处理能力
- HBase 支持低延迟的实时查询
- Iceberg 实现高效的数据湖存储
✅ 最佳实践
- 合理设计数据分层架构
- 实现端到端的 Exactly-Once 语义
- 定期进行表维护和优化
- 建立完善的监控和告警机制
- 做好容错和故障恢复准备
🚀 适用场景
- 实时数据仓库
- 用户行为分析
- IoT 数据处理
- 实时推荐系统
- 金融风控系统
🚫 注意事项:
- 合理控制并行度,避免资源浪费
- 注意数据倾斜问题
- 定期清理过期数据
- 监控各组件的健康状态
通过合理的架构设计和优化,这套技术栈可以支撑 PB 级数据的实时处理需求,为企业的数据驱动决策提供强有力的支持。
上一章 / 下一章
- 上一章:4. Flink CDC.md
- 下一章:6. Flink 监控.md