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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 06:37:01