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

NestJS+Kafka:如何发送已解析为合法JSON的事件值?

问题

我用NestJS服务器作为Kafka主题的事件生产者,已经把Lambda函数订阅到Kafka,会收到如下格式的事件数组。当前payload.value不是标准JSON,而是扁平化的层级键值对结构(类似Java对象的toString输出)。我尝试用正则处理转换成JSON,但代码没法正确解析像"url":"http://"这类片段,而且我不擅长正则表达式。想问有没有办法让Kafka直接发送合法JSON格式的事件值?

事件示例:

[
  {
    "payload": {
        "key": "poll-data",
        "value": "{sensors=[{name=temp_c, id=019428d4-c7b4-759a-90c6-60d66200fa1e}, {name=humidity, id=019428d4-c7b5-777e-96c0-8f2b4024d193}, {name=cloud, id=019428d4-c7b5-777e-96c0-912c8fbbdda8}, {name=wind_degree, id=019428d4-c7b4-759a-90c6-73940dac5464}, {name=wind_kph, id=019428d4-c7b4-759a-90c6-6f436bf2c5af}, {name=wind_dir, id=019428d4-c7b5-777e-96c0-7aadbdb1a480}, {name=pressure_in, id=019428d4-c7b5-777e-96c0-831852ab8f5f}, {name=temp_f, id=019428d4-c7b4-759a-90c6-66c85cf3d48c}, {name=wind_mph, id=019428d4-c7b4-759a-90c6-68d115e0037e}, {name=pressure_mb, id=019428d4-c7b5-777e-96c0-7d11cf6fffb0}, {name=precip_in, id=019428d4-c7b5-777e-96c0-8ad587aa6389}, {name=feelslike_c, id=019428d4-c7b5-777e-96c0-9734b723763c}, {name=feelslike_f, id=019428d4-c7b5-777e-96c0-99928a74a941}, {name=precip_mm, id=019428d4-c7b5-777e-96c0-8787d8c78c18}, {name=windchill_c, id=019428d4-c7b5-777e-96c0-9f65ff08ba97}, {name=windchill_f, id=019428d4-c7b5-777e-96c0-a29cf3a52dfa}, {name=heatindex_f, id=019428d4-c7b5-777e-96c0-a99c1b28cc1e}, {name=heatindex_c, id=019428d4-c7b5-777e-96c0-a47350b244af}, {name=vis_km, id=019428d4-c7b5-777e-96c0-b6e2802cea7a}, {name=vis_miles, id=019428d4-c7b5-777e-96c0-bb3d61493057}, {name=dewpoint_c, id=019428d4-c7b5-777e-96c0-ae2ce16a5b5a}, {name=dewpoint_f, id=019428d4-c7b5-777e-96c0-b0edc620a969}, {name=gust_kph, id=019428d4-c7b5-777e-96c0-c72e86b77b21}, {name=uv, id=019428d4-c7b5-777e-96c0-bfc9f10f8b82}, {name=gust_mph, id=019428d4-c7b5-777e-96c0-c118676f85a4}], connection={serverAddress=http://api.weatherapi.com/v1/current.json, lat=37.201206, long=-3.739167, token=6b48d6699c574b3ab02125435345325325717240612645475675}, deviceId=eb63465a-b263-4b46-9b69-ff13e35fb2d7}",
        "timestamp": 1735852800784,
        "topic": "current-weather-adapter",
        "partition": 0,
        "offset": 8
    }
}
]
解决方案

核心思路是修改NestJS生产者的代码,确保发送前把数据序列化为标准JSON字符串——你现在收到的格式是对象默认toString()方法的输出,并非标准JSON。

1. 手动序列化消息内容

如果是手动构建Kafka消息,直接用JSON.stringify()把数据转换成标准JSON字符串,示例代码如下:

import { Injectable } from '@nestjs/common';
import { ClientKafka } from '@nestjs/microservices';

@Injectable()
export class WeatherProducer {
  constructor(private readonly kafkaClient: ClientKafka) {}

  async sendWeatherData(data: any) {
    // 手动序列化为标准JSON
    const jsonValue = JSON.stringify(data);
    await this.kafkaClient.emit('current-weather-adapter', {
      key: 'poll-data',
      value: jsonValue,
    });
  }
}

2. 配置Kafka客户端的序列化器

如果使用NestJS的自动消息发送机制(比如@MessagePattern),可以在Kafka客户端配置中指定序列化规则,确保消息值被序列化为JSON:

// 模块配置文件
import { Module } from '@nestjs/common';
import { ClientsModule, Transport } from '@nestjs/microservices';

@Module({
  imports: [
    ClientsModule.register([
      {
        name: 'KAFKA_CLIENT',
        transport: Transport.KAFKA,
        options: {
          client: {
            brokers: ['你的Kafka broker地址:9092'],
          },
          consumer: {
            groupId: 'weather-consumer-group',
          },
          // 配置序列化器
          serializer: {
            serialize: (value: any) => JSON.stringify(value),
          },
        },
      },
    ]),
  ],
})
export class AppModule {}

3. 排查中间件干扰

检查NestJS应用中的拦截器、管道或中间件,确认没有代码意外将JSON对象转换成了toString()格式(比如某些日志中间件或自定义数据转换逻辑)。

修改完成后,Kafka发送的payload.value会是标准JSON字符串,Lambda函数可以直接用JSON.parse()解析,无需再用正则处理。

内容的提问来源于stack exchange,提问作者Vladyslav Googlya

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 05:54:53