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
相关产品推荐
相关产品推荐

