Kafka消费者未按预期更新活动状态问题排查
问题描述
我用node-schedule实现了活动状态自动更新——当当前时间和活动开始/结束时间一致时,自动切换状态,这个功能运行正常。但服务器重启后,所有调度任务都会丢失,导致活动状态无法自动更新。
为解决该问题,我尝试将活动信息存入Kafka Topic,期望服务器重启后通过消费这些消息恢复调度,但方案未生效,无法定位问题。
相关代码及数据
生产者代码
export class CampaignProducer { private readonly producer: Producer; private readonly campaignTopic = CAMPAIGN_TOPIC.campaignTopic; constructor(private readonly kafka: Kafka) { this.producer = this.kafka.producer(); } async sendCampaignMessage(campaignId: string, startDate: Date): Promise<void> { await this.producer.connect(); await this.producer.send({ topic: this.campaignTopic, messages: [{ value: JSON.stringify({ campaignId, startDate }) }], }); } }
消费者代码
export class CampaignConsumer implements OnModuleInit { private consumer: Consumer; private readonly campaignTopic = CAMPAIGN_TOPIC.campaignTopic; private readonly groupId = CAMPAIGN_TOPIC.campaignGroup; constructor(private readonly kafka: Kafka, private readonly campaignService: ClientCampaignService) { this.consumer = this.kafka.consumer({ groupId: this.groupId }); } async onModuleInit() { await this.consumeSchedulerInfo(); } async consumeSchedulerInfo() { try { await this.consumer.connect(); await this.consumer.subscribe({ topic: this.campaignTopic, fromBeginning: true }); await this.consumer.run({ autoCommit: true, eachMessage: async ({ message }: EachMessagePayload) => { try { const { campaignId, startDate } = JSON.parse(message.value.toString()); await this.campaignService.scheduleCampaignStart(campaignId, startDate); } catch (error) { console.error('Error processing scheduler message:', error); } }, }); } catch (error) { console.error('Kafka consumer error:', error); } } }
Kafka存储的数据
{"campaignId":"64ff1042a4c66b4c3bc70b05","startDate":"2023-09-11T13:05:00.000Z","endDate":"2023-09-11T13:07:00.000Z"}
MongoDB中的活动数据
{ "_id" : ObjectId("64ff1042a4c66b4c3bc70b05"), "campaignName" : "new campaign", "storeId" : ObjectId("64f02cb62d554f12069d3a36"), "startDateWithTime" : ISODate("2023-09-11T13:05:00.000+0000"), "endDateWithTime" : ISODate("2023-09-11T13:07:00.000+0000"), "campaignStatus" : "Not Started", }
活动创建代码(发送Kafka消息的入口)
async addCampaign(data: Partial<ICampaign>, files) { const directory = path.join(process.cwd(), 'public', UPLOAD_DIRECTORY.CAMPAIGN); await createDirectoryIfNotExists(directory); const filePaths = await uploadFiles(files, directory); const campaignDataWithFiles = { ...data, files: filePaths }; const campaign = await this.campaignModel.create(campaignDataWithFiles); await this.scheduleCampaignAction( campaign._id.toString(), campaign.startDateWithTime, this.updateStatusToInProgress.bind(this) ); await this.scheduleCampaignAction(campaign._id.toString(), campaign.endDateWithTime, this.updateStatusToClosed.bind(this)); return campaign; }
Kafka服务文件代码
export class KafkaService implements OnModuleInit { private kafka: Kafka; private readonly kafkaConfig: KafkaConfig; private campaignProducer: CampaignProducer; private campaignConsumer: CampaignConsumer; constructor(private readonly campaignService: ClientCampaignService) { this.kafkaConfig = { clientId: 'your-client-id', brokers: ['localhost:9092'], }; this.kafka = new Kafka(this.kafkaConfig); this.campaignProducer = new CampaignProducer(this.kafka); this.campaignConsumer = new CampaignConsumer(this.kafka, this.campaignService); } async onModuleInit(): Promise<void> { await this.campaignConsumer.consumeSchedulerInfo(); } getCampaignConsumer(): CampaignConsumer { return this.campaignConsumer; } getCampaignProducer(): CampaignProducer { return this.campaignProducer; } }
问题排查及修复方案
1. 生产者未发送完整调度信息
从Kafka存储的数据看,消息包含endDate,但生产者sendCampaignMessage方法只序列化了campaignId和startDate,且活动结束时间的调度未触发Kafka消息发送。
- 修复:修改生产者方法支持传入
endDate,并在scheduleCampaignAction中针对结束时间调用生产者发送消息。
2. 消费者仅处理开始时间调度
消费者当前只调用scheduleCampaignStart,未处理活动结束时间的调度恢复。
- 修复:修改消费者消息处理逻辑,同时解析
startDate和endDate,分别触发对应的调度方法(如新增scheduleCampaignEnd)。
3. 生产者重复连接Kafka
每次调用sendCampaignMessage都执行producer.connect(),会导致重复连接异常。
- 修复:给生产者添加连接状态判断,仅初始化时连接一次:
export class CampaignProducer { private readonly producer: Producer; private readonly campaignTopic = CAMPAIGN_TOPIC.campaignTopic; private isConnected = false; constructor(private readonly kafka: Kafka) { this.producer = this.kafka.producer(); } async connect(): Promise<void> { if (!this.isConnected) { await this.producer.connect(); this.isConnected = true; } } async sendCampaignMessage(campaignId: string, startDate?: Date, endDate?: Date): Promise<void> { await this.connect(); await this.producer.send({ topic: this.campaignTopic, messages: [{ value: JSON.stringify({ campaignId, startDate, endDate }) }], }); } }
4. 未过滤过期活动的调度恢复
服务器重启后,消费消息时未判断时间是否过期,会创建无效的过期调度任务。
- 修复:在消费者中添加时间判断:
eachMessage: async ({ message }: EachMessagePayload) => { try { const { campaignId, startDate, endDate } = JSON.parse(message.value.toString()); const now = new Date(); if (startDate && new Date(startDate) > now) { await this.campaignService.scheduleCampaignStart(campaignId, startDate); } if (endDate && new Date(endDate) > now) { await this.campaignService.scheduleCampaignEnd(campaignId, endDate); } } catch (error) { console.error('Error processing scheduler message:', error); } }
5. 消费者组偏移量问题
即使设置fromBeginning: true,若消费者组已提交过偏移量,Kafka会从上次位置开始消费。
- 修复:重置消费者组偏移量,或使用新的
groupId测试,确保能消费历史消息。
6. 活动创建时未触发Kafka消息发送
当前scheduleCampaignAction仅创建内存调度任务,未同步发送Kafka消息,导致Kafka中消息来源异常。
- 修复:在
scheduleCampaignAction中添加Kafka消息发送逻辑:async scheduleCampaignAction(campaignId: string, targetDate: Date, action: () => Promise<void>) { schedule.scheduleJob(targetDate, action); if (action === this.updateStatusToInProgress.bind(this)) { await this.kafkaService.getCampaignProducer().sendCampaignMessage(campaignId, targetDate); } else if (action === this.updateStatusToClosed.bind(this)) { await this.kafkaService.getCampaignProducer().sendCampaignMessage(campaignId, undefined, targetDate); } }
内容的提问来源于stack exchange,提问作者John Oliver
相关产品推荐
相关产品推荐

