Kafka多主题消息匹配合并及索引架构优化方案咨询
First, let's align on your context to make sure we're solving the right problem:
Previously, your indexing system relied on Hibernate Search, Lucene, and JBoss HornetQ, following a database-centric model: the app sent metadata (unique DB record identifiers) to HornetQ, and consumers fetched full records (including child objects) from the database using that metadata. Now you're replacing HornetQ with Kafka to eliminate the database-centric dependency. Data creation works smoothly, but edits (e.g., a parent entity with two child objects) pose a challenge:
- When pulling data for user display, you push the full record to Kafka Topic1
- When the user edits and submits only parent-level changes, you push just the parent data to Topic2
- You need to merge the child data from Topic1 with the updated parent data from Topic2 for indexing (since your index only supports delete-then-insert, no direct updates)
Let's break down your questions one by one:
1. How to correlate messages between Topic1 and Topic2? Can we set the same Message ID for both topics?
You should not use Kafka's built-in message IDs/offsets for cross-topic correlation—these are unique per topic partition and aren't designed to be shared across topics. Instead, use a business-level unique identifier (like the parent entity's UUID or primary key) embedded in the payload of both Topic1 and Topic2 messages. For example, add a field like entity_id to every message in both topics. This lets your consumer easily group all messages related to the same parent entity, regardless of which topic they came from.
If you tried to force a custom "message ID" across topics, you'd have to manage ID generation entirely on your end (Kafka doesn't support setting shared IDs), which adds unnecessary complexity. The business entity ID is the natural, reliable correlation key here.
2. Can a single topic solve this problem?
Absolutely—this can simplify your architecture and avoid cross-topic correlation entirely. Here's how to implement it:
- Add an
event_typefield to your message payload (e.g.,FULL_RECORDfor the full data pushed when the user views the entity,PARTIAL_UPDATEfor the parent-only edit submission) - Ensure all messages for the same
entity_idare sent to the same Kafka partition (useentity_idas the partition key). This guarantees messages for a single entity are processed in the order they were sent. - Your consumer can maintain a local cache (or use an external cache like Redis) of full records. When it receives a
PARTIAL_UPDATEmessage, it fetches the cached full record, merges the updated parent data, then uses the merged record to perform the delete-then-insert for indexing. It also updates the cache with the new merged record for future updates.
A single topic cuts down on operational overhead (fewer topics to manage) and removes the need to coordinate between two separate streams.
3. Are there better design options?
Here are a couple of more streamlined approaches to consider:
Option 1: Fetch the full record at submission time
Instead of only pushing parent-level data when the user submits an edit, modify your submission logic to fetch the full parent+child record from the database right after the edit is saved. Then push this complete merged record directly to Kafka. This way, your consumer doesn't need to handle merging at all—they just take the full record from Kafka and run the delete-then-insert. This is the simplest approach if your system can tolerate the extra DB query at submission time.
Option 2: Use Kafka Streams for stream joining
If you need to keep Topic1 and Topic2 separate, use Kafka Streams to build a stream processing job that joins messages from both topics on entity_id. The job listens for new messages in either topic, maintains a state store of full records, and outputs the merged complete record to a new "index-ready" topic. Your index consumer only needs to listen to this single topic, simplifying its logic.
Option 3: Event Sourcing pattern
Model all changes to your entity as events (e.g., EntityCreated, ParentDataUpdated, ChildAdded, etc.) and push these events to Kafka. Your index consumer can then "replay" these events to reconstruct the full current state of the entity whenever it needs to update the index. This is more scalable for complex entities with frequent partial updates, as it avoids storing full records repeatedly.
内容的提问来源于stack exchange,提问作者hithendra sharma

