全部笔记All notes

Flink 结合 Kafka、Protobuf、Cassandra、Redis 构建实时数仓

阅读 7m 34s7m 34s read

概述

在 实时数据仓库 & 流式数据处理 场景下,本文将详细介绍如何使用 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(高性能缓存)

  • 存储内容: 热点数据、实时指标
  • 特点: 内存存储、亚毫秒延迟
  • 使用场景: 实时大屏、监控告警

数据流转过程

  1. 数据采集: 业务系统产生事件数据
  2. 序列化: 使用 Protobuf 序列化为二进制格式
  3. 消息投递: 发送到 Kafka Topic
  4. 流式计算: Flink 消费并实时处理
  5. 结果存储: 并行写入 Cassandra 和 Redis
  6. 数据服务: 提供查询 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.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);
    }
}
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());
        // 实际场景可以发送到监控系统、发送邮件等
    }
}
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 压缩减少网络传输
  • 并行度: 根据数据量设置合适的并行度
  • 状态后端: 使用 RocksDB 处理大状态
  • 检查点: 合理设置检查点间隔
  • 内存管理: 调整 TaskManager 内存配置

3. Cassandra 优化

  • 批量写入: 使用批处理减少写入次数
  • 分区键设计: 避免热点分区
  • 压缩策略: 使用时间窗口压缩
  • TTL: 设置数据过期时间

4. Redis 优化

  • Pipeline: 批量操作减少网络往返
  • 连接池: 复用连接避免频繁创建
  • 过期策略: 合理设置 key 过期时间
  • 持久化: 根据需求选择 RDB 或 AOF

常见问题

Q1: Protobuf 反序列化失败怎么办?

解决方案:

  1. 检查 proto 文件版本是否一致
  2. 添加异常处理和死信队列
  3. 记录失败数据便于排查
@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; // 跳过该消息
    }
}

解决方案:

  1. 启用检查点并配置状态后端
  2. 设置重启策略
env.setStateBackend(new RocksDBStateBackend("hdfs://namenode:9000/flink/checkpoints"));
env.setRestartStrategy(RestartStrategies.fixedDelayRestart(3, Time.seconds(10)));

Q3: Cassandra 写入延迟高?

解决方案:

  1. 使用异步写入
  2. 调整一致性级别
  3. 批量写入优化
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 内存占用过高?

解决方案:

  1. 设置合理的过期时间
  2. 使用 Redis 集群分片
  3. 定期清理冷数据
  4. 监控内存使用情况

总结

本文详细介绍了如何使用 Flink + Kafka + Protobuf + Cassandra + Redis 构建实时数据仓库:

核心要点

  1. 架构设计: 各组件职责明确,相互配合
  2. 数据流转: 从采集到存储的完整链路
  3. 实时计算: 窗口聚合、实时告警等功能
  4. 双写模式: Cassandra 存历史,Redis 存热点
  5. 性能优化: 各组件的调优要点

技术栈总结

组件版本作用
Kafka2.8.0+消息队列,Protobuf 数据传输
Flink1.16.0+流式计算引擎
Protobuf3.19+高效序列化协议
Cassandra3.11+分布式存储,历史数据
Redis6.0+内存缓存,实时查询

适用场景

  • 实时用户行为分析
  • 实时监控告警
  • 实时推荐系统
  • 实时风控系统
  • IoT 数据处理

通过这套架构,我们可以构建一个高效、可靠、可扩展的实时数据处理平台!

相关文章

  • Flink DataStream API 基础
  • Flink 窗口操作详解
  • Flink 状态管理
  • Kafka 基础教程
  • Cassandra 入门指南
  • Redis 最佳实践