基于Kafka的事件驱动架构中如何执行逻辑操作?——以电商用户与商品微服务场景为例
在事件驱动架构中实现商品创建的验证逻辑(避免数据副本)
你面临的问题其实是事件驱动架构中强一致性验证的典型场景——既要摆脱同步RPC的耦合,又不想依赖本地副本带来的数据不一致风险。下面结合你的电商场景,给出两种实用的实现方案,都基于Kafka这类事件流平台:
方案一:基于Kafka请求-响应模式的同步验证(强一致)
这种方式本质是用事件流模拟"同步查询",但不直接依赖RPC,而是通过Kafka的请求/响应主题实现服务间的异步通信,同时保证验证的实时性,而且不需要在商品服务存储用户套餐副本。
实现步骤:
- 定义事件主题:
- 创建
user-plan-requests主题:商品服务发送查询用户套餐类型的请求事件 - 创建
user-plan-responses主题:用户服务返回查询结果的响应事件
- 创建
- 商品服务的创建逻辑:
// 商品服务代码 const { v4: uuidv4 } = require('uuid'); async function createProduct(userId, newProduct) { // 1. 先查询本地商品数量 const productsCount = await query("SELECT COUNT(*) FROM products WHERE user_id = ?", [userId]); if (productsCount >= 10) { // 2. 数量达标,需要验证用户是否付费 const requestId = uuidv4(); // 生成唯一请求ID用于关联响应 // 发送查询用户套餐的请求事件到Kafka await kafkaProducer.send({ topic: 'user-plan-requests', messages: [{ key: userId, value: JSON.stringify({ requestId, userId }) }] }); // 3. 监听对应的响应事件(设置超时时间,避免无限等待) let response; try { response = await new Promise((resolve, reject) => { const timeout = setTimeout(() => reject(new Error('User plan query timeout')), 5000); const consumer = kafkaConsumer.subscribe({ topic: 'user-plan-responses', fromBeginning: false }); consumer.run({ eachMessage: async ({ message }) => { const data = JSON.parse(message.value.toString()); if (data.requestId === requestId) { clearTimeout(timeout); await consumer.stop(); resolve(data.isPaid); } } }); }); } catch (err) { return { error: 'Failed to verify user plan: ' + err.message }; } // 4. 根据响应判断是否允许创建 if (!response) { return { error: 'Cannot create new product: free plan limit reached' }; } } // 5. 允许创建,插入商品后发送事件到Kafka const productId = await query("INSERT INTO products (user_id, name, ...) VALUES (?, ?, ...)", [userId, newProduct.name, ...]); await kafkaProducer.send({ topic: 'product-created', messages: [{ key: userId, value: JSON.stringify({ userId, productId, product: newProduct }) }] }); return { success: true, productId }; } - 用户服务的处理逻辑:
// 用户服务代码 kafkaConsumer.subscribe({ topic: 'user-plan-requests' }); kafkaConsumer.run({ eachMessage: async ({ message }) => { const { requestId, userId } = JSON.parse(message.value.toString()); // 查询用户是否为付费套餐 const [user] = await query("SELECT is_paid FROM users WHERE id = ?", [userId]); const isPaid = user?.is_paid || false; // 发送响应事件 await kafkaProducer.send({ topic: 'user-plan-responses', messages: [{ key: userId, value: JSON.stringify({ requestId, userId, isPaid }) }] }); } });
优势:
- 不需要在商品服务存储用户套餐数据,彻底避免数据同步问题
- 基于事件流通信,服务间耦合度远低于RPC(可更换服务实现,只要事件格式兼容)
- 保证验证的强一致性,适合对数据准确性要求极高的场景
注意点:
- 需要处理请求超时、重复响应等异常情况(比如用请求ID做去重处理)
- 属于"异步等待"模式,比纯同步RPC稍复杂,但灵活性和可扩展性更强
方案二:异步验证+补偿机制(最终一致)
如果你的业务可以接受短暂的"不一致窗口"(比如用户创建商品后,几秒内被撤销),那么可以采用先创建、后验证、不符合则补偿的方式,完全异步,不需要等待用户服务的响应。
实现步骤:
- 定义事件主题:
product-created:商品创建成功后发送的事件product-creation-invalidated:验证不通过时发送的撤销事件
- 商品服务的创建逻辑:
// 商品服务代码 async function createProduct(userId, newProduct) { // 1. 直接插入商品(不做严格前置验证) const productId = await query("INSERT INTO products (user_id, name, ...) VALUES (?, ?, ...)", [userId, newProduct.name, ...]); // 2. 发送商品创建事件 await kafkaProducer.send({ topic: 'product-created', messages: [{ key: userId, value: JSON.stringify({ userId, productId, product: newProduct }) }] }); return { success: true, productId, note: 'Your product is being verified, it will be visible shortly' }; } - 新增验证服务(或由用户服务承担):
这个服务专门监听product-created事件,负责验证用户是否符合创建条件:// 验证服务代码 kafkaConsumer.subscribe({ topic: 'product-created' }); kafkaConsumer.run({ eachMessage: async ({ message }) => { const { userId, productId } = JSON.parse(message.value.toString()); // 1. 查询用户套餐类型 const [user] = await query("SELECT is_paid FROM users WHERE id = ?", [userId]); const isPaid = user?.is_paid || false; // 2. 查询该用户的商品总数 const [countResult] = await query("SELECT COUNT(*) as total FROM products WHERE user_id = ?", [userId]); const productsCount = countResult.total; // 3. 如果不符合条件,触发撤销逻辑 if (productsCount > 10 && !isPaid) { // 删除违规商品 await query("DELETE FROM products WHERE id = ?", [productId]); // 发送撤销事件,通知其他服务(比如用户通知服务) await kafkaProducer.send({ topic: 'product-creation-invalidated', messages: [{ key: userId, value: JSON.stringify({ userId, productId, reason: 'Free plan product limit exceeded' }) }] }); } } });
优势:
- 完全异步,商品服务响应速度极快,不需要等待任何外部服务
- 服务间完全解耦,商品服务不需要知晓验证逻辑的存在
- 适合高并发场景,或对响应速度要求高于即时一致性的业务
注意点:
- 用户可能会短暂看到自己创建的商品,之后被删除,需要在前端做友好提示(比如"正在验证,稍后可见")
- 需要处理验证服务的故障情况(比如用死信队列存储处理失败的事件,后续重试)
两种方案的选择建议
- 如果你的业务要求即时强一致性(绝对不能让用户创建不符合条件的商品),优先选方案一
- 如果你的业务可以接受最终一致性,且追求高并发、低耦合,选方案二
内容的提问来源于stack exchange,提问作者Ari Seyhun
相关产品推荐
相关产品推荐

