概述
在 实时数据仓库 & 流式数据处理 场景下,本文将详细介绍如何使用 Flink + Kafka + Protobuf + Cassandra + Redis 来构建一个高吞吐、低延迟、可扩展的流式数据处理系统。
核心特性
- 高吞吐: Kafka 作为消息队列,支持百万级 TPS
- 低延迟: Flink 毫秒级流处理,Redis 亚毫秒级查询
- 强一致性: Flink 的 exactly-once 语义保证
- 高可用: 所有组件都支持分布式部署
- 灵活扩展: 水平扩展能力强
架构设计
整体架构图
┌─────────────┐ ┌─────────────┐ ┌─────────────┐
│ Producer │────▶│ Kafka │────▶│ Flink │
└─────────────┘ └─────────────┘ └──────┬──────┘
(Protobuf格式) │
┌────────┴────────┐
│ │
▼ ▼
┌─────────────┐ ┌─────────────┐
│ Cassandra │ │ Redis │
└─────────────┘ └─────────────┘
(历史数据存储) (热点数据缓存)
组件职责
1. Kafka(消息队列)
- 作用: 数据输入源,消息缓冲
- 数据格式: Protobuf 二进制协议
- 特点: 高吞吐、持久化、分区机制
2. Flink(流式计算引擎)
- 核心功能:
- 从 Kafka 消费 Protobuf 数据
- 实时数据清洗和转换
- 复杂事件处理(CEP)
- 窗口聚合计算
- 输出: 双写 Cassandra 和 Redis
3. Cassandra(分布式数据库)
- 存储内容: 历史数据、聚合结果
- 特点: 高可用、线性扩展、时序数据友好
- 使用场景: 数据回溯、历史查询
4. Redis(高性能缓存)
- 存储内容: 热点数据、实时指标
- 特点: 内存存储、亚毫秒延迟
- 使用场景: 实时大屏、监控告警
数据流转过程
- 数据采集: 业务系统产生事件数据
- 序列化: 使用 Protobuf 序列化为二进制格式
- 消息投递: 发送到 Kafka Topic
- 流式计算: Flink 消费并实时处理
- 结果存储: 并行写入 Cassandra 和 Redis
- 数据服务: 提供查询 API 供下游使用
环境准备
版本要求
- Flink: 1.16.0+
- Kafka: 2.8.0+
- Cassandra: 3.11+
- Redis: 6.0+
- Java: 8 或 11
- Protobuf: 3.19+
Maven 依赖
<dependencies>
<!-- Flink 核心依赖 -->
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-streaming-java</artifactId>
<version>1.16.0</version>
</dependency>
<!-- Kafka Connector -->
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-connector-kafka</artifactId>
<version>1.16.0</version>
</dependency>
<!-- Protobuf -->
<dependency>
<groupId>com.google.protobuf</groupId>
<artifactId>protobuf-java</artifactId>
<version>3.19.4</version>
</dependency>
<!-- Cassandra Connector -->
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-connector-cassandra</artifactId>
<version>1.16.0</version>
</dependency>
<!-- Redis Connector -->
<dependency>
<groupId>redis.clients</groupId>
<artifactId>jedis</artifactId>
<version>4.2.3</version>
</dependency>
</dependencies>
实战步骤
1. Kafka 数据源配置
Kafka 作为流数据的入口,消息格式使用 Protobuf 以提高传输效率。
1.1 定义 Protobuf 数据结构
创建 event.proto 文件定义数据结构:
syntax = "proto3";
package com.example.flink.proto;
option java_package = "com.example.flink.proto";
option java_outer_classname = "EventProto";
// 事件消息定义
message Event {
string user_id = 1; // 用户ID
string event_type = 2; // 事件类型
int64 timestamp = 3; // 时间戳
map<string, string> properties = 4; // 扩展属性
// 嵌套消息示例
message Location {
double latitude = 1;
double longitude = 2;
}
Location location = 5; // 位置信息(可选)
}
// 批量事件消息
message EventBatch {
repeated Event events = 1;
}
编译 Protobuf 生成 Java 类:
# 安装 protoc 编译器
# macOS: brew install protobuf
# Linux: apt-get install protobuf-compiler
# 编译生成 Java 代码
protoc --java_out=src/main/java src/main/proto/event.proto
1.2 Kafka Topic 创建
# 创建 Topic(3个分区,2个副本)
kafka-topics.sh --create \
--bootstrap-server localhost:9092 \
--topic user-events \
--partitions 3 \
--replication-factor 2 \
--config retention.ms=86400000 \
--config compression.type=lz4
1.3 生产 Kafka 消息
创建 Kafka Producer 发送 Protobuf 消息:
import com.example.flink.proto.EventProto.Event;
import org.apache.kafka.clients.producer.*;
import java.util.Properties;
import java.util.concurrent.TimeUnit;
public class EventProducer {
private final KafkaProducer<String, byte[]> producer;
private final String topic;
public EventProducer(String bootstrapServers, String topic) {
Properties props = new Properties();
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG,
"org.apache.kafka.common.serialization.StringSerializer");
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG,
"org.apache.kafka.common.serialization.ByteArraySerializer");
// 性能优化配置
props.put(ProducerConfig.ACKS_CONFIG, "1");
props.put(ProducerConfig.BATCH_SIZE_CONFIG, 16384);
props.put(ProducerConfig.LINGER_MS_CONFIG, 10);
props.put(ProducerConfig.COMPRESSION_TYPE_CONFIG, "lz4");
this.producer = new KafkaProducer<>(props);
this.topic = topic;
}
public void sendEvent(String userId, String eventType) {
Event event = Event.newBuilder()
.setUserId(userId)
.setEventType(eventType)
.setTimestamp(System.currentTimeMillis())
.putProperties("source", "web")
.putProperties("version", "1.0")
.build();
ProducerRecord<String, byte[]> record =
new ProducerRecord<>(topic, userId, event.toByteArray());
producer.send(record, new Callback() {
@Override
public void onCompletion(RecordMetadata metadata, Exception e) {
if (e != null) {
System.err.println("发送失败: " + e.getMessage());
} else {
System.out.printf("发送成功: topic=%s, partition=%d, offset=%d\n",
metadata.topic(), metadata.partition(), metadata.offset());
}
}
});
}
public void close() {
producer.close(5, TimeUnit.SECONDS);
}
public static void main(String[] args) {
EventProducer producer = new EventProducer("localhost:9092", "user-events");
// 模拟发送事件
for (int i = 0; i < 100; i++) {
producer.sendEvent("user_" + i, "click");
try {
Thread.sleep(100);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
producer.close();
}
}
2. Flink 流处理实现
2.1 自定义 Protobuf 反序列化器
实现 Kafka 的反序列化接口来解析 Protobuf 数据:
import com.example.flink.proto.EventProto.Event;
import org.apache.flink.api.common.serialization.DeserializationSchema;
import org.apache.flink.api.common.typeinfo.TypeInformation;
import org.apache.flink.connector.kafka.source.reader.deserializer.KafkaRecordDeserializationSchema;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import java.io.IOException;
public class ProtobufDeserializationSchema implements DeserializationSchema<Event> {
@Override
public Event deserialize(byte[] message) throws IOException {
if (message == null) {
return null;
}
try {
return Event.parseFrom(message);
} catch (Exception e) {
throw new IOException("Failed to deserialize Protobuf message", e);
}
}
@Override
public boolean isEndOfStream(Event nextElement) {
return false;
}
@Override
public TypeInformation<Event> getProducedType() {
return TypeInformation.of(Event.class);
}
}
2.2 Flink 主程序实现
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.functions.timestamps.BoundedOutOfOrdernessTimestampExtractor;
import org.apache.flink.streaming.api.windowing.time.Time;
import org.apache.flink.connector.kafka.source.KafkaSource;
import org.apache.flink.connector.kafka.source.enumerator.initializer.OffsetsInitializer;
import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import java.time.Duration;
public class RealtimeWarehouseJob {
public static void main(String[] args) throws Exception {
// 1. 创建执行环境
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// 设置并行度和检查点
env.setParallelism(3);
env.enableCheckpointing(60000); // 1分钟检查点
env.getCheckpointConfig().setMinPauseBetweenCheckpoints(30000);
env.getCheckpointConfig().setCheckpointTimeout(60000);
env.getCheckpointConfig().setMaxConcurrentCheckpoints(1);
// 2. 配置 Kafka Source(Flink 1.14+ 新API)
KafkaSource<Event> kafkaSource = KafkaSource.<Event>builder()
.setBootstrapServers("localhost:9092")
.setTopics("user-events")
.setGroupId("flink-consumer-group")
.setStartingOffsets(OffsetsInitializer.latest())
.setValueOnlyDeserializer(new ProtobufDeserializationSchema())
.build();
// 3. 创建数据流并设置水印策略
DataStream<Event> eventStream = env
.fromSource(kafkaSource,
WatermarkStrategy.<Event>forBoundedOutOfOrderness(Duration.ofSeconds(5))
.withTimestampAssigner((event, timestamp) -> event.getTimestamp()),
"Kafka Source")
.name("Event Stream");
// 4. 数据转换和处理(见下一节)
DataStream<UserEventStats> statsStream = processEvents(eventStream);
// 5. 双写输出
// 写入 Cassandra
statsStream.addSink(new CassandraSink())
.name("Cassandra Sink");
// 写入 Redis
statsStream.addSink(new RedisSink())
.name("Redis Sink");
// 6. 启动作业
env.execute("Realtime Data Warehouse Job");
}
private static DataStream<UserEventStats> processEvents(DataStream<Event> eventStream) {
// 具体处理逻辑见下一节
return null;
}
}
3. 实时计算逻辑
3.1 定义聚合结果类
import java.io.Serializable;
public class UserEventStats implements Serializable {
private String userId;
private String eventType;
private long eventCount;
private long lastUpdateTime;
private long windowStart;
private long windowEnd;
// 构造函数、getter/setter 省略
@Override
public String toString() {
return String.format("UserEventStats{userId='%s', eventType='%s', count=%d, window=[%d,%d]}",
userId, eventType, eventCount, windowStart, windowEnd);
}
}
3.2 实现多维度统计
import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows;
import org.apache.flink.streaming.api.windowing.time.Time;
import org.apache.flink.streaming.api.functions.windowing.ProcessWindowFunction;
import org.apache.flink.streaming.api.windowing.windows.TimeWindow;
import org.apache.flink.util.Collector;
private static DataStream<UserEventStats> processEvents(DataStream<Event> eventStream) {
// 1. 按用户和事件类型分组,进行窗口聚合
DataStream<UserEventStats> userStats = eventStream
.keyBy(event -> event.getUserId() + "_" + event.getEventType())
.window(TumblingEventTimeWindows.of(Time.minutes(5))) // 5分钟滚动窗口
.process(new ProcessWindowFunction<Event, UserEventStats, String, TimeWindow>() {
@Override
public void process(String key, Context context,
Iterable<Event> events,
Collector<UserEventStats> out) {
long count = 0;
String userId = null;
String eventType = null;
for (Event event : events) {
if (userId == null) {
userId = event.getUserId();
eventType = event.getEventType();
}
count++;
}
UserEventStats stats = new UserEventStats();
stats.setUserId(userId);
stats.setEventType(eventType);
stats.setEventCount(count);
stats.setWindowStart(context.window().getStart());
stats.setWindowEnd(context.window().getEnd());
stats.setLastUpdateTime(System.currentTimeMillis());
out.collect(stats);
}
})
.name("User Event Statistics");
// 2. 实时告警:检测异常行为
DataStream<UserEventStats> alerts = userStats
.filter(stats -> stats.getEventCount() > 100) // 5分钟内事件超过100次
.name("Anomaly Detection");
alerts.addSink(new AlertSink()).name("Alert Sink");
return userStats;
}
// 告警输出
public static class AlertSink extends RichSinkFunction<UserEventStats> {
@Override
public void invoke(UserEventStats value, Context context) {
System.err.printf("[ALERT] 用户 %s 在5分钟内产生了 %d 次 %s 事件\n",
value.getUserId(), value.getEventCount(), value.getEventType());
// 实际场景可以发送到监控系统、发送邮件等
}
}
3.3 使用 Flink SQL 实现
import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;
import org.apache.flink.table.api.Table;
public static void processBySQL(StreamExecutionEnvironment env,
DataStream<Event> eventStream) {
StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);
// 注册数据流为表
tableEnv.createTemporaryView("events", eventStream,
$"userId", $"eventType", $"timestamp".rowtime(), $"properties");
// SQL 查询:5分钟窗口统计
String sql = """
SELECT
userId,
eventType,
COUNT(*) as eventCount,
TUMBLE_START(timestamp, INTERVAL '5' MINUTE) as windowStart,
TUMBLE_END(timestamp, INTERVAL '5' MINUTE) as windowEnd
FROM events
GROUP BY
userId,
eventType,
TUMBLE(timestamp, INTERVAL '5' MINUTE)
""";
Table resultTable = tableEnv.sqlQuery(sql);
// 转换回 DataStream
DataStream<UserEventStats> resultStream = tableEnv
.toDataStream(resultTable, UserEventStats.class);
}
4. Cassandra 持久化存储
4.1 创建 Cassandra 表结构
-- 创建 Keyspace
CREATE KEYSPACE IF NOT EXISTS flink_warehouse
WITH replication = {'class': 'SimpleStrategy', 'replication_factor': 3};
USE flink_warehouse;
-- 用户事件统计表(按时间分区)
CREATE TABLE IF NOT EXISTS user_event_stats (
user_id TEXT,
event_type TEXT,
window_date DATE,
window_time TIMESTAMP,
event_count BIGINT,
last_update_time TIMESTAMP,
PRIMARY KEY ((user_id, window_date), window_time, event_type)
) WITH CLUSTERING ORDER BY (window_time DESC, event_type ASC)
AND default_time_to_live = 2592000 -- 30天过期
AND compaction = {'class': 'TimeWindowCompactionStrategy'};
-- 创建索引便于查询
CREATE INDEX IF NOT EXISTS idx_event_type ON user_event_stats (event_type);
-- 用户画像表
CREATE TABLE IF NOT EXISTS user_profiles (
user_id TEXT PRIMARY KEY,
total_events BIGINT,
last_active_time TIMESTAMP,
favorite_event_types MAP<TEXT, BIGINT>,
daily_active_hours LIST<INT>,
tags SET<TEXT>
);
4.2 实现 Cassandra Sink
import com.datastax.driver.core.*;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.streaming.api.functions.sink.RichSinkFunction;
import java.time.Instant;
import java.time.LocalDate;
import java.time.ZoneId;
public class CassandraSink extends RichSinkFunction<UserEventStats> {
private transient Cluster cluster;
private transient Session session;
private transient PreparedStatement insertStatement;
private transient PreparedStatement updateProfileStatement;
@Override
public void open(Configuration parameters) throws Exception {
// 连接配置
cluster = Cluster.builder()
.addContactPoint("localhost")
.withPort(9042)
.withCredentials("cassandra", "cassandra")
.withPoolingOptions(new PoolingOptions()
.setConnectionsPerHost(HostDistance.LOCAL, 2, 4)
.setMaxRequestsPerConnection(HostDistance.LOCAL, 32768))
.build();
session = cluster.connect("flink_warehouse");
// 准备语句
String insertCql = """
INSERT INTO user_event_stats
(user_id, event_type, window_date, window_time, event_count, last_update_time)
VALUES (?, ?, ?, ?, ?, ?)
""";
insertStatement = session.prepare(insertCql);
String updateProfileCql = """
UPDATE user_profiles
SET total_events = total_events + ?,
last_active_time = ?,
favorite_event_types[?] = ?
WHERE user_id = ?
""";
updateProfileStatement = session.prepare(updateProfileCql);
}
@Override
public void invoke(UserEventStats stats, Context context) throws Exception {
try {
// 1. 插入事件统计数据
LocalDate windowDate = Instant.ofEpochMilli(stats.getWindowStart())
.atZone(ZoneId.systemDefault())
.toLocalDate();
BoundStatement boundStatement = insertStatement.bind(
stats.getUserId(),
stats.getEventType(),
windowDate,
new Date(stats.getWindowStart()),
stats.getEventCount(),
new Date(stats.getLastUpdateTime())
);
session.execute(boundStatement);
// 2. 更新用户画像
BoundStatement profileStatement = updateProfileStatement.bind(
stats.getEventCount(),
new Date(stats.getLastUpdateTime()),
stats.getEventType(),
stats.getEventCount(),
stats.getUserId()
);
session.execute(profileStatement);
} catch (Exception e) {
System.err.println("写入 Cassandra 失败: " + e.getMessage());
throw e;
}
}
@Override
public void close() throws Exception {
if (session != null) {
session.close();
}
if (cluster != null) {
cluster.close();
}
}
}
5. Redis 缓存实现
5.1 Redis 数据结构设计
# 1. 实时统计数据
# Hash: 存储用户最新统计信息
HSET user:stats:{userId} event_count 100 last_update_time 1234567890
# 2. 排行榜
# Sorted Set: 活跃用户排行榜
ZADD active_users:{date} 100 user_123
# 3. 时间序列数据
# Redis TimeSeries (需要安装 RedisTimeSeries 模块)
TS.ADD user:ts:{userId}:click 1234567890 10
# 4. 实时告警
# List: 告警队列
LPUSH alerts:queue "{userId:123,type:anomaly,time:1234567890}"
5.2 高性能 Redis Sink 实现
import redis.clients.jedis.*;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.streaming.api.functions.sink.RichSinkFunction;
import java.util.HashMap;
import java.util.Map;
public class RedisSink extends RichSinkFunction<UserEventStats> {
private transient JedisPool jedisPool;
private final String redisHost;
private final int redisPort;
private final String password;
private final int dbIndex;
public RedisSink(String host, int port, String password, int dbIndex) {
this.redisHost = host;
this.redisPort = port;
this.password = password;
this.dbIndex = dbIndex;
}
@Override
public void open(Configuration parameters) throws Exception {
JedisPoolConfig poolConfig = new JedisPoolConfig();
poolConfig.setMaxTotal(100);
poolConfig.setMaxIdle(10);
poolConfig.setMinIdle(5);
poolConfig.setTestOnBorrow(true);
poolConfig.setTestOnReturn(true);
poolConfig.setTestWhileIdle(true);
if (password != null && !password.isEmpty()) {
jedisPool = new JedisPool(poolConfig, redisHost, redisPort,
2000, password, dbIndex);
} else {
jedisPool = new JedisPool(poolConfig, redisHost, redisPort,
2000, null, dbIndex);
}
}
@Override
public void invoke(UserEventStats stats, Context context) throws Exception {
try (Jedis jedis = jedisPool.getResource()) {
// 使用 Pipeline 批量操作提高性能
Pipeline pipeline = jedis.pipelined();
// 1. 更新用户统计信息
String userKey = "user:stats:" + stats.getUserId();
Map<String, String> userStats = new HashMap<>();
userStats.put("event_type", stats.getEventType());
userStats.put("event_count", String.valueOf(stats.getEventCount()));
userStats.put("last_update_time", String.valueOf(stats.getLastUpdateTime()));
userStats.put("window_start", String.valueOf(stats.getWindowStart()));
userStats.put("window_end", String.valueOf(stats.getWindowEnd()));
pipeline.hmset(userKey, userStats);
pipeline.expire(userKey, 3600); // 1小时过期
// 2. 更新活跃用户排行榜
String rankKey = "active_users:" + getCurrentDateKey();
pipeline.zadd(rankKey, stats.getEventCount(), stats.getUserId());
pipeline.expire(rankKey, 86400); // 1天过期
// 3. 更新事件类型统计
String eventTypeKey = "event_type:" + stats.getEventType() + ":" + getCurrentHourKey();
pipeline.hincrBy(eventTypeKey, "count", stats.getEventCount());
pipeline.hincrBy(eventTypeKey, "users", 1);
pipeline.expire(eventTypeKey, 7200); // 2小时过期
// 4. 发布实时事件(用于实时大屏)
String channel = "realtime:events";
String message = String.format("{\"userId\":\"%s\",\"eventType\":\"%s\",\"count\":%d}",
stats.getUserId(), stats.getEventType(), stats.getEventCount());
pipeline.publish(channel, message);
// 执行批量操作
pipeline.sync();
} catch (Exception e) {
System.err.println("写入 Redis 失败: " + e.getMessage());
throw e;
}
}
@Override
public void close() throws Exception {
if (jedisPool != null) {
jedisPool.close();
}
}
private String getCurrentDateKey() {
return java.time.LocalDate.now().toString();
}
private String getCurrentHourKey() {
return java.time.LocalDateTime.now().format(
java.time.format.DateTimeFormatter.ofPattern("yyyy-MM-dd-HH")
);
}
}
5.3 Redis 查询服务
import org.springframework.web.bind.annotation.*;
import redis.clients.jedis.Jedis;
import java.util.*;
@RestController
@RequestMapping("/api/realtime")
public class RealtimeQueryController {
private final JedisPool jedisPool;
// 获取用户实时统计
@GetMapping("/user/{userId}/stats")
public Map<String, Object> getUserStats(@PathVariable String userId) {
try (Jedis jedis = jedisPool.getResource()) {
Map<String, String> stats = jedis.hgetAll("user:stats:" + userId);
return new HashMap<>(stats);
}
}
// 获取活跃用户排行榜
@GetMapping("/ranking/active-users")
public List<Map<String, Object>> getActiveUserRanking(
@RequestParam(defaultValue = "10") int limit) {
try (Jedis jedis = jedisPool.getResource()) {
String key = "active_users:" + getCurrentDateKey();
Set<Tuple> ranking = jedis.zrevrangeWithScores(key, 0, limit - 1);
List<Map<String, Object>> result = new ArrayList<>();
int rank = 1;
for (Tuple tuple : ranking) {
Map<String, Object> item = new HashMap<>();
item.put("rank", rank++);
item.put("userId", tuple.getElement());
item.put("eventCount", tuple.getScore());
result.add(item);
}
return result;
}
}
}
完整示例
项目结构
realtime-warehouse/
├── src/main/java/
│ ├── proto/ # Protobuf 定义
│ ├── producer/ # Kafka 生产者
│ ├── flink/ # Flink 作业
│ │ ├── deserializer/
│ │ ├── function/
│ │ ├── sink/
│ │ └── job/
│ └── api/ # 查询接口
├── src/main/resources/
│ ├── application.yaml
│ └── log4j2.xml
├── docker-compose.yml # 本地环境
└── pom.xml
Docker Compose 本地环境
version: '3.8'
services:
zookeeper:
image: confluentinc/cp-zookeeper:7.3.0
environment:
ZOOKEEPER_CLIENT_PORT: 2181
ZOOKEEPER_TICK_TIME: 2000
kafka:
image: confluentinc/cp-kafka:7.3.0
depends_on:
- zookeeper
ports:
- "9092:9092"
environment:
KAFKA_BROKER_ID: 1
KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
cassandra:
image: cassandra:3.11
ports:
- "9042:9042"
environment:
- CASSANDRA_CLUSTER_NAME=flink-cluster
- CASSANDRA_DC=dc1
- CASSANDRA_ENDPOINT_SNITCH=GossipingPropertyFileSnitch
redis:
image: redis:7-alpine
ports:
- "6379:6379"
command: redis-server --requirepass redis123
redis-insight:
image: redislabs/redisinsight:latest
ports:
- "8001:8001"
性能优化
1. Kafka 优化
- 分区策略: 按 userId hash 分区,保证同一用户数据在同一分区
- 批量发送: 设置 batch.size 和 linger.ms
- 压缩: 使用 LZ4 压缩减少网络传输
2. Flink 优化
- 并行度: 根据数据量设置合适的并行度
- 状态后端: 使用 RocksDB 处理大状态
- 检查点: 合理设置检查点间隔
- 内存管理: 调整 TaskManager 内存配置
3. Cassandra 优化
- 批量写入: 使用批处理减少写入次数
- 分区键设计: 避免热点分区
- 压缩策略: 使用时间窗口压缩
- TTL: 设置数据过期时间
4. Redis 优化
- Pipeline: 批量操作减少网络往返
- 连接池: 复用连接避免频繁创建
- 过期策略: 合理设置 key 过期时间
- 持久化: 根据需求选择 RDB 或 AOF
常见问题
Q1: Protobuf 反序列化失败怎么办?
解决方案:
- 检查 proto 文件版本是否一致
- 添加异常处理和死信队列
- 记录失败数据便于排查
@Override
public Event deserialize(byte[] message) throws IOException {
try {
return Event.parseFrom(message);
} catch (InvalidProtocolBufferException e) {
// 发送到死信队列
deadLetterProducer.send(new ProducerRecord<>("dead-letter-topic", message));
return null; // 跳过该消息
}
}
Q2: Flink 作业重启后状态丢失?
解决方案:
- 启用检查点并配置状态后端
- 设置重启策略
env.setStateBackend(new RocksDBStateBackend("hdfs://namenode:9000/flink/checkpoints"));
env.setRestartStrategy(RestartStrategies.fixedDelayRestart(3, Time.seconds(10)));
Q3: Cassandra 写入延迟高?
解决方案:
- 使用异步写入
- 调整一致性级别
- 批量写入优化
ResultSetFuture future = session.executeAsync(boundStatement);
Futures.addCallback(future, new FutureCallback<ResultSet>() {
@Override
public void onSuccess(ResultSet result) {
// 成功处理
}
@Override
public void onFailure(Throwable t) {
// 失败重试
}
});
Q4: Redis 内存占用过高?
解决方案:
- 设置合理的过期时间
- 使用 Redis 集群分片
- 定期清理冷数据
- 监控内存使用情况
总结
本文详细介绍了如何使用 Flink + Kafka + Protobuf + Cassandra + Redis 构建实时数据仓库:
核心要点
- 架构设计: 各组件职责明确,相互配合
- 数据流转: 从采集到存储的完整链路
- 实时计算: 窗口聚合、实时告警等功能
- 双写模式: Cassandra 存历史,Redis 存热点
- 性能优化: 各组件的调优要点
技术栈总结
| 组件 | 版本 | 作用 |
|---|---|---|
| Kafka | 2.8.0+ | 消息队列,Protobuf 数据传输 |
| Flink | 1.16.0+ | 流式计算引擎 |
| Protobuf | 3.19+ | 高效序列化协议 |
| Cassandra | 3.11+ | 分布式存储,历史数据 |
| Redis | 6.0+ | 内存缓存,实时查询 |
适用场景
- 实时用户行为分析
- 实时监控告警
- 实时推荐系统
- 实时风控系统
- IoT 数据处理
通过这套架构,我们可以构建一个高效、可靠、可扩展的实时数据处理平台!
相关文章
- Flink DataStream API 基础
- Flink 窗口操作详解
- Flink 状态管理
- Kafka 基础教程
- Cassandra 入门指南
- Redis 最佳实践