Flink 源码与运行原理

流式计算 / 源码阅读 / LESSON 04

State & StateBackend 深度解析

从 State API 到 StateBackend,理解状态如何被组织、存储和在恢复时重建。

阅读时间
90 分钟
学习路径
Flink 源码与运行原理
内容来源
Flink 深度笔记

预计阅读时间: 90 分钟 前置阅读: doc-03(Checkpoint) 下一次阅读: doc-05(Window)


1. 状态分类与内部表示

1.1 三种状态类型

类型作用域存储后端典型场景内部结构
KeyedStatekeyBy() 后的每个 KeyStateBackend (Heap/RocksDB)GROUP BY, 窗口聚合, 计数KeyGroup → {Key → Value}
OperatorState整个算子实例StateBackendKafka Source offsetList<StateHandle>
BroadcastState广播到所有 SubtaskStateBackend动态配置分发Key → Value (每个 Subtask 全量)

1.2 KeyedState 的物理组织 — KeyGroup

Key 不直接映射到 Subtask,而是经过两层映射:

Key → hash(key) % maxParallelism → KeyGroup → KeyGroupRangeAssignment → Subtask

例如 maxParallelism=128, parallelism=4:
  KeyGroup 分配:
    Subtask 0: KeyGroup [0, 31]
    Subtask 1: KeyGroup [32, 63]
    Subtask 2: KeyGroup [64, 95]
    Subtask 3: KeyGroup [96, 127]

Key "user_123" → hash=42 → KeyGroup=42 → Subtask 1 (KeyGroup [32,63])
// KeyGroupRangeAssignment.java — 核心分配算法
public static int computeKeyGroupForKeyHash(int keyHash, int maxParallelism) {
    return MathUtils.murmurHash(keyHash) % maxParallelism;
}

public static int computeOperatorIndexForKeyGroup(
        int maxParallelism, int parallelism, int keyGroup) {
    // 将 KeyGroup 均匀分配到 Subtask
    return keyGroup * parallelism / maxParallelism;
}

为什么需要 KeyGroup 这一层?

  • KeyGroup 数量(maxParallelism)创建后固定不变
  • 当并行度改变时,只需重新分配 KeyGroup→Subtask 映射
  • Key→KeyGroup 映射永远不变 → Rescale 时状态可以恢复

1.3 五种 KeyedState 类型

// KeyedStateBackend.java — 状态类型
public interface KeyedStateBackend<K> extends KeyedStateFactory {
    // 1. ValueState<T>: 单值状态
    <T> ValueState<T> getPartitionedState(
        ValueStateDescriptor<T> descriptor);

    // 2. ListState<T>: 列表状态
    <T> ListState<T> getPartitionedState(
        ListStateDescriptor<T> descriptor);

    // 3. MapState<K,V>: Map 状态 (支持迭代)
    <T, UK, UV> MapState<UK, UV> getPartitionedState(
        MapStateDescriptor<UK, UV> descriptor);

    // 4. ReducingState<T>: 归约状态 (自动 reduce)
    <T> ReducingState<T> getPartitionedState(
        ReducingStateDescriptor<T> descriptor);

    // 5. AggregatingState<IN, OUT>: 聚合状态
    <IN, ACC, OUT> AggregatingState<IN, OUT> getPartitionedState(
        AggregatingStateDescriptor<IN, ACC, OUT> descriptor);
}

ValueState vs MapState 的选择:

维度ValueStateMapState
存储单值Key-Value 集合
State TTL支持每个 Entry 独立 TTL
迭代不支持支持 keys()/values()/entries() 迭代
RocksDB 存储单 CF Key-Value多个 CF Key-Value (带前缀)
适用简单计数、累加窗口内元素、GROUP BY 多个 group

2. StateBackend 深度对比

2.1 三种 Backend 对比

维度HashMapStateBackendEmbeddedRocksDBStateBackendChangelogStateBackend
存储位置Java HeapRocksDB (本地磁盘)DFS Changelog+本地 Materialized
状态大小上限~100MB (受 GC 限制)TB 级 (受本地磁盘限制)中(受 DFS 限制)
访问速度极快(直接堆访问)快(有序列化开销,~1-5μs)快(本地读)+ DFS 写
序列化仅在 Checkpoint 时每次读写都需要写入时序列化到 DFS
Checkpoint 格式全量 (Java 对象 → 字节)增量 (SST 文件 diff)Changelog 日志
恢复速度快(纯内存加载)中(需下载 SST + 重放 WAL)极快(Apply Changelog)
GC 影响大(状态在 Heap 中)无(状态在 Off-Heap/磁盘)
适用场景小状态 + 低延迟大状态 + 稳定吞吐需要快速恢复的场景

2.2 StateBackend 接口

// StateBackend.java — 创建 KeyedStateBackend 工厂
public interface StateBackend extends java.io.Serializable {
    // 创建 KeyedStateBackend(keyBy 后的状态)
    <K> KeyedStateBackend<K> createKeyedStateBackend(
        Environment env, JobID jobID, String operatorIdentifier,
        TypeSerializer<K> keySerializer, int numberOfKeyGroups,
        KeyGroupRange keyGroupRange, TaskKvStateRegistry kvStateRegistry,
        TtlTimeProvider ttlTimeProvider, MetricGroup metricGroup,
        Collection<KeyedStateHandle> stateHandles, // 恢复时的状态
        CloseableRegistry cancelStreamRegistry,
        double managedMemoryFraction) throws Exception;

    // 创建 OperatorStateBackend(Kafka Source offset 等)
    OperatorStateBackend createOperatorStateBackend(
        Environment env, String operatorIdentifier,
        Collection<OperatorStateHandle> stateHandles,
        CloseableRegistry cancelStreamRegistry) throws Exception;
}

3. EmbeddedRocksDBStateBackend 内部原理

3.1 架构总览

TaskManager JVM
│
├── RocksDBKeyedStateBackend
│   │
│   ├── ColumnFamily per State:
│   │   ├── CF "windowState"    (MapState 的每个 key-value 用一个 RocksDB key)
│   │   ├── CF "countState"     (ValueState 单值)
│   │   └── CF "default"
│   │
│   ├── WriteBatch (批量写入,减少 RocksDB 调用)
│   │   └── 积累多次 state update → 批量 flush
│   │
│   ├── RocksDB 实例(C++ JNI 调用)
│   │   ├── MemTable (write buffer) → Immutable MemTable → flush
│   │   ├── SST Files: L0 (overlapping) → L1 → ... → L6 (non-overlapping)
│   │   ├── Block Cache (读缓存,在 Managed Memory 中)
│   │   └── WAL (write-ahead log, 进程崩溃恢复)
│   │
│   └── Managed Memory 分配:
│       ├── writeBufferSize * maxWriteBufferNumber → MemTable
│       ├── blockCacheSize → Block Cache
│       └── 总 managed memory 在 taskmanager.memory.managed.size 中配置

3.2 RocksDB 中的 Key 编码

RocksDB 是 KV 存储,Flink 通过组合编码将 KeyedState 存储在 RocksDB 中:

RocksDB Key 格式:
┌──────────────────────────────────────────────────┐
│ KeyGroup(2B) │ Key(Variable) │ Namespace(Variable) │
└──────────────────────────────────────────────────┘

例如:
  KeyGroup=42, Key="user_123", Namespace=Window[10:00,10:01)
  → RocksDB Key = [0x002A][user_123][window-10:00-10:01]

为什么要加 KeyGroup 前缀?
  - 恢复时按 KeyGroup 扫描(不需要全量扫描)
  - KeyGroup 是恢复和 Rescale 的最小单位

3.3 RocksDB 写路径

// RocksDBWriteBatchWrapper.java — 批量写入
public class RocksDBWriteBatchWrapper implements AutoCloseable {
    private final WriteBatch writeBatch;
    private final WriteOptions writeOptions;  // disable WAL 等配置

    public void put(ColumnFamilyHandle cf, byte[] key, byte[] value) throws RocksDBException {
        if (writeBatch.count() >= batchSize || sizeInBytes >= batchSizeInBytes) {
            flush();        // 批量满 → 写入 RocksDB
        }
        writeBatch.put(cf, key, value);
        sizeInBytes += key.length + value.length;
    }

    private void flush() throws RocksDBException {
        rocksDB.write(writeOptions, writeBatch);
        writeBatch.clear();
    }
}

写入优化

  1. WriteBatch 批量积累:默认 batchSize=2MB,减少 RocksDB write() 调用
  2. 禁用 WAL(可选):Checkpoint 提供容错,RocksDB WAL 可禁用 → 写性能大幅提升
  3. MemTable 快速路径:大部分写入直接写 MemTable(内存),异步 flush 到磁盘

3.4 RocksDB 读路径

// RocksDBValueState.java — 读取 ValueState
public V value() throws IOException {
    byte[] key = serializeKeyAndNamespace(currentKey, namespace);
    byte[] valueBytes = backend.db.get(
        currentKeyGroup,   // KeyGroup 范围查询
        columnFamily,      // State 对应的 CF
        key
    );

    if (valueBytes == null) {
        return defaultValue;
    }
    return serializer.deserialize(valueBytes);
}

读的层次(RocksDB 内部):

MemTable → Immutable MemTable → Block Cache → SST Files (L0→L6)
                                     ↑
                              Managed Memory 中的 Block Cache
                              缓存频繁读取的 SST Block

RocksDB 的序列化开销:每次 get()put() 都需要将 Java 对象序列化为字节数组。这是 RocksDB Backend 的主要性能开销来源。对于访问非常频繁的小状态,HashMapStateBackend 没有这个开销。


4. HashMapStateBackend — 纯堆内存状态

4.1 内部结构

// HeapKeyedStateBackend.java — 核心存储
public class HeapKeyedStateBackend<K> extends AbstractKeyedStateBackend<K> {
    // 存储: Map<StateName, StateTable>
    // StateTable = Map<KeyGroup, Map<K, State>>

    private final Map<String, StateTable<K, ?, ?>> registeredKVStates;

    // 每个 StateDescriptor 对应一个 StateTable
    // StateTable 内部是一组 CopyOnWriteStateMap
    // (每个 KeyGroup 一个 CopyOnWriteStateMap)
}

4.2 CopyOnWriteStateMap — 并发安全的状态读写

HeapKeyedStateBackend 的并发模型:

每个 KeyGroup 有独立的 CopyOnWriteStateMap<K, N, S>
  - 读写分离: 读取无锁,写入时拷贝整个 Map
  - 适用于小状态的场景(小于 100MB)
  - GC 压力: 每次 Checkpoint 生成的快照是堆上的新对象

CopyOnWriteStateMap 结构:
  StateMapEntry[] table (hash table)
  └── StateMapEntry { key, namespace, state, next(in-chain) }

4.3 Checkpoint 快照

// HeapKeyedStateBackend.snapshot() — 全量快照
public void snapshot(long checkpointId, long timestamp,
        CheckpointStreamFactory streamFactory,
        CheckpointOptions checkpointOptions) throws Exception {

    // 1. 遍历所有 State → 序列化为字节数组
    for (StateTable stateTable : registeredKVStates.values()) {
        // 2. 按 KeyGroup 分组序列化
        for (int keyGroup : keyGroupRange) {
            // 3. 序列化 KeyGroup 内的所有 Key→State
            DataOutputView outputView = ...;
            stateTable.writeStateInKeyGroup(outputView, keyGroup);
        }
    }
    // 4. 写入 DFS (HDFS/S3)
    // 全量写入 → 状态越大,Checkpoint 越慢
}

HashMapStateBackend 的 Checkpoint 限制

  • 全量 Checkpoint,每次都要序列化和上传全部状态
  • 状态 > 1GB 时,Checkpoint 耗时可能超过 1 分钟
  • 状态 > 100MB 时,GC 压力显著增加
  • 推荐:状态 < 100MB → HashMapStateBackend;> 100MB → RocksDB

5. ChangelogStateBackend — 基于日志的状态

5.1 设计思路

ChangelogStateBackend 将状态变更记录到分布式日志中:

┌─────────────────────────────────────────────┐
│          ChangelogStateBackend              │
│                                             │
│  状态写入:  state.update(value)             │
│      │                                      │
│      ├──→ 写入本地 Materialized Store (快)   │
│      └──→ 追加 Changelog Log (DFS)          │
│                                             │
│  恢复时:                                    │
│      └──→ 从 DFS 读取 Changelog             │
│           → Apply 到本地 Store              │
│           → 恢复完成(无需下载全量状态)      │
└─────────────────────────────────────────────┘

5.2 Changelog vs 传统 Checkpoint

维度传统 CheckpointChangelog Backend
恢复时间下载所有 SST → 重放 WALApply Changelog 日志
恢复数据量O(状态大小) 如 100GBO(Changelog 增量) 如 100MB
写入开销只在 Checkpoint 时写 DFS每次状态更新都写 DFS(可能 durably 或 periodically)
一致性Checkpoint 间隔内无持久化Periodic materialize + continuous log

6. State TTL — 状态过期机制

6.1 TTL 配置

StateTtlConfig ttlConfig = StateTtlConfig
    .newBuilder(Time.hours(24))                     // TTL 24 小时
    .setUpdateType(UpdateType.OnCreateAndWrite)      // 写时更新时间戳
    .setStateVisibility(StateVisibility.NeverReturnExpired)  // 不返回过期值
    .setTtlTimeCharacteristic(TtlTimeCharacteristic.ProcessingTime)  // 处理时间
    .cleanupFullSnapshot()                           // Checkpoint 时全量清理
    .cleanupIncrementally(10, false)                 // 每次访问时增量清理
    .cleanupInRocksdbCompactFilter()                 // RocksDB Compaction 时清理
    .build();

6.2 TTL 清理策略详解

策略机制开销适用
cleanupFullSnapshotCheckpoint 时扫描所有状态,丢弃过期项Checkpoint 开销增加默认策略,最安全
cleanupIncrementally每次 state.value() 访问时检查并清理过期项每次访问有微小开销访问频繁但状态量大的场景
cleanupInRocksdbCompactFilterRocksDB Compaction 时通过 CompactionFilter 清理零读取开销大状态 + RocksDB + 不要求精确到期
cleanupFullSnapshot + cleanupIncrementally组合均衡生产推荐

6.3 TTL 的内部实现

// TtlValueStateWrapper.java — TTL 包装器
class TtlValueStateWrapper<T> extends AbstractTtlState<T> implements ValueState<T> {
    // 每次写时记录 timestamp
    @Override
    public void update(T value) {
        long timestamp = timeProvider.currentTimestamp();  // 当前处理时间/事件时间
        wrappedState.update(TtlValue.wrap(value, timestamp));
    }

    // 每次读时检查 TTL
    @Override
    public T value() {
        TtlValue<T> ttlValue = wrappedState.value();
        if (ttlValue == null) return null;

        if (expired(ttlValue)) {     // timestamp + ttl < currentTime
            // 根据 visibility 决定是否返回
            if (visibility == NeverReturnExpired) {
                return null;         // 或 defaultValue
            }
        }
        return ttlValue.getValue();
    }
}

TTL 的关键时序问题

  • ProcessingTime: 精确(系统时间定长),但重启后时间重置
  • EventTime: 受 Watermark 推进影响,如果 Watermark 停滞,TTL 不生效 → 状态可能堆积
  • setUpdateType(OnReadAndWrite): 每次读也更新 TTL,防止"热"数据意外过期

7. Managed Memory 与状态的关系

7.1 Managed Memory 在 RocksDB Backend 中的作用

taskmanager.memory.managed.size (默认: flink.size * 0.4)

┌─────────────────────────────────────────────────┐
│            Managed Memory = 1GB                 │
├─────────────────────────────────────────────────┤
│ RocksDB Write Buffer (MemTable)          = 512MB │
│ RocksDB Block Cache (Read Cache)         = 256MB │
│ Sort-Merge Join 缓冲区                    = 128MB │
│ Hash Join Build Side                     = 96MB  │
│ Python UDF 内存                           = 32MB  │
└─────────────────────────────────────────────────┘

RocksDB 核心配置

ROCKSDD 配置默认说明
write_buffer_size64MB单个 MemTable 的大小
max_write_buffer_number2MemTable 的数量(1 active + 1 immutable)
block_cache_size8MB (per CF)读缓存的 Block Cache 大小
target_file_size_base64MBL1 层 SST 文件目标大小
max_bytes_for_level_base256MBL1 层总大小上限

7.2 MemoryManager 源码

// MemoryManager.java — Managed Memory 的池化管理
public class MemoryManager {
    private final long memorySize;           // 总内存大小
    private final int pageSize;              // 内存页大小 (32KB)
    private final int numberOfPages;         // 总页数
    private final ArrayDeque<MemorySegment> freeSegments;  // 空闲段池

    // 分配内存段
    public List<MemorySegment> allocatePages(
            Object owner, int numPages) throws MemoryAllocationException {
        synchronized (lock) {
            if (freeSegments.size() < numPages) {
                throw new MemoryAllocationException(...);
            }
            List<MemorySegment> allocated = new ArrayList<>(numPages);
            for (int i = 0; i < numPages; i++) {
                allocated.add(freeSegments.poll());
            }
            allocatedSegments.put(owner, allocated);  // 记账
            return allocated;
        }
    }

    // 释放内存段
    public void release(List<MemorySegment> segments) {
        synchronized (lock) {
            freeSegments.addAll(segments);
            // 更新记账
        }
    }
}

8. RocksDB 增量 Checkpoint 深度流程

8.1 完整流程

RocksDBStateBackend.snapshotState(checkpointId=42)
  │
  ├── Step 1: flushMemTable()
  │   └── 将 MemTable 刷到 SST 文件 → 生成新的 SST 文件
  │
  ├── Step 2: rocksDB.getLiveFiles()
  │   └── 获取当前所有 SST 文件的列表
  │       eg: [SST_01, SST_02, SST_03, SST_04, SST_05]
  │
  ├── Step 3: chooseSstFilesForUpload()
  │   └── 对比 lastCheckpoint (Chk#41) 的 SST 列表
  │       Chk#41 SST files: [SST_01, SST_02, SST_03]
  │       Current SST files: [SST_01, SST_02, SST_03, SST_04, SST_05]
  │       新增: [SST_04, SST_05] ← 只需上传这两个
  │       删除: [] ← compaction 没有删除文件
  │
  ├── Step 4: uploadFilesToCheckpointFs()
  │   └── 上传 SST_04, SST_05 到 DFS
  │       异步执行: 不影响数据处理
  │
  ├── Step 5: SharedStateRegistry.registerReference()
  │   └── 注册 SST_01, SST_02, SST_03 被 Chk#42 引用
  │       refCount[SST_01] = 2 (Chk#41 + Chk#42)
  │       refCount[SST_02] = 2
  │       refCount[SST_03] = 2
  │       refCount[SST_04] = 1 (新)
  │       refCount[SST_05] = 1 (新)
  │
  └── Step 6: 构造 IncrementalRemoteKeyedStateHandle
      → 返回给 CheckpointCoordinator
      → 包含: uploadedHandles + sharedStateRegistry references

8.2 SharedStateRegistry 的引用计数 GC

SharedStateRegistry 维护:
  Map<SstFileKey, Set<CheckpointId>>  // 每个 SST 文件被哪些 Checkpoint 引用

GC 逻辑:
  Chk#41 过期 (被新 Checkpoint 淘汰或达到 maxCheckpoints):
    → 遍历 Chk#41 的 SST 文件
    → refCount[SST_01]-- (2→1)
    → refCount[SST_02]-- (2→1)
    → refCount[SST_03]-- (2→1)
    → SST_01, SST_02, SST_03 仍被 Chk#42 引用 → 保留

  Chk#42 过期:
    → refCount[SST_01]-- (1→0) → 删除 DFS 上的 SST_01
    → refCount[SST_02]-- (1→0) → 删除 DFS 上的 SST_02
    → ...

9. 状态序列化

9.1 TypeSerializer 体系

TypeSerializer<T>  (Flink 的序列化抽象)
├── BasicTypeSerializer (Int, Long, String, Double, Boolean...)
├── PojoSerializer (POJO 类型)
├── KryoSerializer (通用后备方案)
├── RowSerializer (Row 类型, Table API 内部)
├── TupleSerializer (Tuple 类型)
└── CompositeTypeSerializer (自定义复合类型)

序列化性能对比:
  Kryo (通用)     ~100 ns/object (慢, 但通用)
  PojoSerializer  ~50 ns/object  (中)
  BasicType       ~5 ns/object   (快, 仅基本类型)
  Row             ~20 ns/object  (快, Table API 优化)

9.2 RocksDB 中的序列化路径

State 写入:
  Java Object → TypeSerializer.serialize() → byte[] → RocksDB.put(keyBytes, valueBytes)

State 读取:
  RocksDB.get(key) → byte[] → TypeSerializer.deserialize() → Java Object

Key 编码:
  (KeyGroup, Key, Namespace) → CompositeKeySerialization.serializeKeyGroupNamespaceKey()

  特殊优化: 当 Key 和 Namespace 是固定长度类型时(如 Long),使用高效的编码路径

减少序列化开销的技巧

  1. 使用 POJO 而非 Kryo(env.getConfig().enableForceKryo() → 关闭)
  2. 避免大对象作为 State Value(拆成多个小 ValueState)
  3. RocksDB 的 Block Cache 可以缓存序列化后的字节,减少 deserialize 次数

10. 源码导航(完整版)

文件关键类/方法职责
flink-runtime/.../state/StateBackend.javacreateKeyedStateBackend()StateBackend 接口
flink-runtime/.../state/KeyedStateBackend.javasetCurrentKey(), getPartitionedState()KeyedState Backend
flink-runtime/.../state/heap/HeapKeyedStateBackend.javasnapshot(), notifyCheckpointComplete()堆状态后端
flink-runtime/.../state/heap/CopyOnWriteStateMap.javaget(), put(), remove()堆状态存储结构
flink-state-backends/flink-statebackend-rocksdb/.../RocksDBKeyedStateBackend.javasnapshot()RocksDB 状态后端
flink-state-backends/flink-statebackend-rocksdb/.../RocksDBWriteBatchWrapper.javaput(), flush()RocksDB 批量写入
flink-state-backends/flink-statebackend-rocksdb/.../RocksDBIncrementalCheckpointUtils.javachooseSstFilesForUpload()增量 SST 选择
flink-runtime/.../state/ttl/TtlStateFactory.javacreateValueState(), ...TTL 状态包装器
flink-runtime/.../state/ttl/TtlValueState.javavalue(), update()TTL ValueState 实现
flink-runtime/.../memory/MemoryManager.javaallocatePages(), release()Managed Memory 管理
flink-runtime/.../state/KeyGroupRangeAssignment.javacomputeKeyGroupForHash(), computeOperatorIndexForKeyGroup()KeyGroup 分配
flink-runtime/.../state/SharedStateRegistry.javaregisterReference(), unregisterUnusedState()SST 文件引用计数

11. 常见问题 / 面试题

Q1: RocksDB vs HashMap StateBackend 的选择标准?

A: 以 100MB 为界。以下是详细决策树:

  • 状态 < 100MB → HashMapStateBackend(零序列化开销,低延迟)
  • 状态 > 100MB → EmbeddedRocksDBStateBackend(无 GC 压力,增量 Checkpoint)
  • 状态巨大但需要极快恢复 → ChangelogStateBackend
  • 不确定 → 从 RocksDB 开始,因为它不增加 GC 压力且可增量 Checkpoint

Q2: RocksDB 的 WAL 和 Flink Checkpoint 的关系是什么?

A: 它们互补但不重复:

  • RocksDB WAL:进程级崩溃恢复(TM 重启但本地磁盘还在),恢复 ~秒级
  • Flink Checkpoint:全局容错(TM 挂掉、磁盘损坏、集群故障),恢复 ~分钟级
  • 默认情况下 Flink Checkpoint 前会 flush MemTable,因此 WAL 在 Checkpoint 间的作用有限
  • 如果禁用 WAL(state.backend.rocksdb.wal.disabled: true),依赖 Checkpoint 恢复

Q3: 状态大小如何监控?

A: 关键 Metrics:

  • State Size:每个算子的状态大小(RocksDB 的 estimate-num-keys, size-all-mem-tables, block-cache-usage
  • Checkpoint Size:每次 Checkpoint 的增量/全量大小
  • RocksDB 的 numBytesInMemTable, numBytesInBlockCache, numBytesInCompaction 在 Metrics 中可查

Q4: State TTL 为什么不使用事件时间?

A: 可以使用,但有风险:

  • 如果 Watermark 不推进(Idle Source / 巨大乱序),TTL 永远不会触发
  • 状态会无限堆积,最终 OOM
  • 建议:设置 cleanupIncrementally + cleanupInRocksdbCompactFilter 兜底,即使 TTL 不触发也能通过访问时/Compaction 时清理

Q5: 为什么 RocksDB 状态下做了 KeyBy 后还需要序列化?

A: RocksDB 是 C++ 库(JNI 调用),数据必须从 JVM Heap 拷贝到 C++ 内存。这个过程:

  1. Java Object → TypeSerializer.serialize() → byte[]
  2. byte[] → JNI → C++ RocksDB → MemTable
  3. 读取时:C++ RocksDB → JNI → byte[] → TypeSerializer.deserialize() → Java Object

Q6: OperatorState 和 KeyedState 的区别?举一个具体的例子。

A: 以 Kafka Source 为例:

  • OperatorState(Kafka Source):存储 (partition, offset) 映射。因为 Kafka Partition 的分配不依赖 Key,而是依赖 Subtask。
    • Subtask 0 的 OperatorState: {partition-0: offset=1000, partition-1: offset=2000}
    • Subtask 1 的 OperatorState: {partition-2: offset=1500, partition-3: offset=1800}
  • KeyedState(Window 算子):存储每个 Key 的窗口聚合结果。
    • Key "user_001" 的 KeyedState: {window[10:00]: count=5, window[10:01]: count=3}
    • 这个状态按 Key 所属的 KeyGroup 分布在不同的 Subtask 中

下一步