You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何为Apache Flink开发原始键控状态?自定义布隆过滤器状态问询

Got it, I totally get why you're looking for a more efficient way here—storing the entire BloomFilter instance in ValueState means you have to read/write the whole object every time you update it, which gets pretty costly especially with large filters or remote state backends. The raw state API in Flink is indeed under-documented, so let's break this down step by step, including both a custom raw state implementation and a simpler alternative that might fit your needs.

Why ValueState<BloomFilter> is inefficient

Every time you add an element to the bloom filter, you have to:

  1. Fetch the entire BloomFilter instance from state
  2. Modify it in memory
  3. Serialize the whole instance and write it back to state

This full round-trip is wasteful for bloom filters, since we only need to flip individual bits rather than touch the entire object. Raw state lets us work directly with byte arrays, cutting out the overhead of serializing/deserializing the entire BloomFilter instance.

Option 1: Implement Custom Raw State

Raw state is Flink's lowest-level state API, giving you full control over how state is stored and manipulated. Here's how to build a custom raw state for your bloom filter:

Step 1: Define the Raw State Class

This class handles the bloom filter logic directly on a byte buffer, avoiding object serialization overhead:

public class BloomFilterRawState implements State {
    private final ByteBuffer bitmap;
    private final int numHashFunctions;
    private final int bitSize;

    public BloomFilterRawState(int bitSize, int numHashFunctions) {
        this.bitSize = bitSize;
        this.numHashFunctions = numHashFunctions;
        this.bitmap = ByteBuffer.allocate(bitSize / 8);
    }

    // Add element by flipping relevant bits directly in the byte buffer
    public void add(byte[] element) {
        for (int i = 0; i < numHashFunctions; i++) {
            int hash = Math.abs(MurmurHash3.hash(element, i)) % bitSize;
            int byteIndex = hash / 8;
            int bitIndex = hash % 8;
            bitmap.put(byteIndex, (byte) (bitmap.get(byteIndex) | (1 << bitIndex)));
        }
    }

    // Check if element might exist
    public boolean mightContain(byte[] element) {
        for (int i = 0; i < numHashFunctions; i++) {
            int hash = Math.abs(MurmurHash3.hash(element, i)) % bitSize;
            int byteIndex = hash / 8;
            int bitIndex = hash % 8;
            if ((bitmap.get(byteIndex) & (1 << bitIndex)) == 0) {
                return false;
            }
        }
        return true;
    }

    @Override
    public void clear() {
        bitmap.clear();
        Arrays.fill(bitmap.array(), (byte) 0);
    }

    // Serialize state to byte array for checkpointing
    public byte[] serialize() {
        return bitmap.array();
    }

    // Deserialize state from byte array (for recovery)
    public static BloomFilterRawState deserialize(byte[] data, int bitSize, int numHashFunctions) {
        BloomFilterRawState state = new BloomFilterRawState(bitSize, numHashFunctions);
        state.bitmap.put(data);
        state.bitmap.flip();
        return state;
    }
}

Step 2: Create a State Descriptor

This tells Flink how to serialize/deserialize your custom state:

public class BloomFilterStateDescriptor extends StateDescriptor<BloomFilterRawState, byte[]> {
    private final int bitSize;
    private final int numHashFunctions;

    public BloomFilterStateDescriptor(String name, int bitSize, int numHashFunctions) {
        super(name, TypeInformation.of(byte[].class), null);
        this.bitSize = bitSize;
        this.numHashFunctions = numHashFunctions;
    }

    @Override
    public StateSerializer<BloomFilterRawState> createSerializer(ExecutionConfig executionConfig) {
        return new StateSerializer<BloomFilterRawState>() {
            @Override
            public byte[] serialize(BloomFilterRawState state) throws IOException {
                return state.serialize();
            }

            @Override
            public BloomFilterRawState deserialize(byte[] data) throws IOException {
                return BloomFilterRawState.deserialize(data, bitSize, numHashFunctions);
            }

            @Override
            public BloomFilterRawState copy(BloomFilterRawState from) throws IOException {
                BloomFilterRawState copy = new BloomFilterRawState(bitSize, numHashFunctions);
                copy.bitmap.put(from.bitmap.array());
                copy.bitmap.flip();
                return copy;
            }
        };
    }
}

Step 3: Use the Custom State in an Operator

Now you can use this state in your process function:

public class BloomFilterProcessFunction extends RichProcessFunction<String, String> {
    private BloomFilterRawState bloomFilterState;

    @Override
    public void open(Configuration parameters) throws Exception {
        super.open(parameters);
        // Configure bloom filter: 1MB bitmap, 5 hash functions
        BloomFilterStateDescriptor descriptor = new BloomFilterStateDescriptor(
                "custom-bloom-filter",
                1024 * 1024 * 8,
                5
        );
        bloomFilterState = getRuntimeContext().getState(descriptor);
    }

    @Override
    public void processElement(String value, Context ctx, Collector<String> out) throws Exception {
        byte[] elementBytes = value.getBytes(StandardCharsets.UTF_8);
        if (!bloomFilterState.mightContain(elementBytes)) {
            bloomFilterState.add(elementBytes);
            out.collect(value); // Output only new elements
        }
    }
}

Option 2: Simpler Alternative - ValueState<byte[]>

If you don't want to dive into raw state, you can store just the bloom filter's bitmap as a byte[] in ValueState. This avoids serializing the entire BloomFilter object, which is already a big win:

// Scala implementation example
class BloomFilterFunction extends RichProcessFunction[String, String] {
  private var bitmapState: ValueState[Array[Byte]] = _
  private val bitSize = 1024 * 1024 * 8 // 1MB bitmap
  private val numHashFunctions = 5

  override def open(parameters: Configuration): Unit = {
    val descriptor = new ValueStateDescriptor[Array[Byte]](
      "bloom-filter-bitmap",
      createTypeInformation[Array[Byte]]
    )
    bitmapState = getRuntimeContext.getState(descriptor)
    // Initialize bitmap if state is empty
    if (bitmapState.value() == null) {
      bitmapState.update(new Array[Byte](bitSize / 8))
    }
  }

  override def processElement(value: String, ctx: Context, out: Collector[String]): Unit = {
    val bytes = value.getBytes(StandardCharsets.UTF_8)
    var elementExists = true
    val bitmap = bitmapState.value()

    for (i <- 0 until numHashFunctions) {
      val hash = Math.abs(MurmurHash3.hash(bytes, i)) % bitSize
      val byteIdx = hash / 8
      val bitIdx = hash % 8

      if ((bitmap(byteIdx) & (1 << bitIdx)) == 0) {
        elementExists = false
        // Flip the bit to 1
        bitmap(byteIdx) = (bitmap(byteIdx) | (1 << bitIdx)).toByte
      }
    }

    if (!elementExists) {
      bitmapState.update(bitmap)
      out.collect(value)
    }
  }
}

Key Notes

  • Compatibility: When using raw state, you're responsible for handling schema changes (e.g., if you change the bit size or hash function count, old checkpoints won't be compatible).
  • Dynamic Scaling: If you need a bloom filter that can grow, raw state will require more complex logic to handle bitmap resizing while maintaining consistency.
  • Memory: For very large bloom filters, consider using Flink's managed memory instead of heap memory to avoid GC pressure—this requires working with off-heap byte buffers.

内容的提问来源于stack exchange,提问作者Moein Hosseini

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.15 08:24:27