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

Kafka新手求教:如何发送含随机数与时间戳的多值消息?

在Kafka中发送包含多个值的消息:标准实现方式

嘿,刚上手Kafka的话,这个问题确实很常见——毕竟默认的ProducerRecord只能传递一个value字段,但要把随机数和时间戳这类多字段数据塞进消息里,其实有几种非常标准的实现方式,我给你拆解下:

方案1:打包成JSON字符串(最通用、易上手)

这是业界最常用的方式:把多个字段(比如你的随机数+时间戳)封装成一个JSON对象,转成字符串后作为ProducerRecord的value发送。这种方式的好处是跨语言兼容,不管消费者用Java、Python还是其他语言,都能轻松解析JSON拿到字段。

修改后的代码示例

首先你需要引入JSON序列化的依赖,比如Jackson(如果用Maven的话,在pom.xml里加):

<dependency>
    <groupId>com.fasterxml.jackson.core</groupId>
    <artifactId>jackson-databind</artifactId>
    <version>2.15.2</version> <!-- 用最新稳定版即可 -->
</dependency>

然后调整你的生产者代码,注意把Random和SimpleDateFormat移到循环外(避免重复创建对象浪费资源):

import com.fasterxml.jackson.databind.ObjectMapper;
import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.ProducerRecord;

import java.text.SimpleDateFormat;
import java.util.Date;
import java.util.HashMap;
import java.util.Map;
import java.util.Properties;
import java.util.Random;

public class MultiFieldKafkaProducer {
    public static void main(String[] args) throws InterruptedException {
        // 1. 配置Kafka生产者(根据你的集群信息调整)
        Properties props = new Properties();
        props.put("bootstrap.servers", "localhost:9092");
        props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
        props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");

        KafkaProducer<String, String> producer = new KafkaProducer<>(props);

        // 2. 把重复创建的对象移到循环外,提升性能
        Random rand = new Random();
        SimpleDateFormat formatter = new SimpleDateFormat("dd/MM/yyyy HH:mm:ss");
        ObjectMapper objectMapper = new ObjectMapper();

        try {
            while (true) {
                for (int key = 0; key < 10000; key++) {
                    int randomNumber = rand.nextInt(50) + 1;
                    String timestamp = formatter.format(new Date());

                    // 3. 把多字段打包成Map,再转成JSON字符串
                    Map<String, Object> messageContent = new HashMap<>();
                    messageContent.put("randomNumber", randomNumber);
                    messageContent.put("timestamp", timestamp);
                    String jsonValue = objectMapper.writeValueAsString(messageContent);

                    // 4. 发送包含JSON的消息
                    ProducerRecord<String, String> record = new ProducerRecord<>("java-topic", Integer.toString(1), jsonValue);
                    producer.send(record);

                    Thread.sleep(10000);
                }
            }
        } catch (Exception e) {
            e.printStackTrace();
        } finally {
            producer.close(); // 建议在程序退出时关闭生产者,避免资源泄漏
        }
    }
}

方案2:自定义Java对象+自定义序列化器(强类型场景)

如果你的项目是纯Java生态,并且希望用强类型的方式处理消息(避免JSON解析的潜在错误),可以定义一个包含所需字段的Java类,然后自定义Kafka的序列化器来序列化这个对象。

步骤示例

  1. 定义自定义消息类:
public class CustomMessage {
    private int randomNumber;
    private String timestamp;

    // 构造方法、Getter、Setter
    public CustomMessage(int randomNumber, String timestamp) {
        this.randomNumber = randomNumber;
        this.timestamp = timestamp;
    }

    public int getRandomNumber() { return randomNumber; }
    public void setRandomNumber(int randomNumber) { this.randomNumber = randomNumber; }
    public String getTimestamp() { return timestamp; }
    public void setTimestamp(String timestamp) { this.timestamp = timestamp; }
}
  1. 自定义序列化器:
import com.fasterxml.jackson.databind.ObjectMapper;
import org.apache.kafka.common.serialization.Serializer;

import java.util.Map;

public class CustomMessageSerializer implements Serializer<CustomMessage> {
    private final ObjectMapper objectMapper = new ObjectMapper();

    @Override
    public void configure(Map<String, ?> configs, boolean isKey) {
        // 可以添加自定义配置,比如序列化规则
    }

    @Override
    public byte[] serialize(String topic, CustomMessage data) {
        try {
            return objectMapper.writeValueAsBytes(data);
        } catch (Exception e) {
            throw new RuntimeException("Failed to serialize CustomMessage", e);
        }
    }

    @Override
    public void close() {
        // 关闭资源(如果需要)
    }
}
  1. 配置生产者并发送消息:
// 调整生产者配置,指定自定义的value序列化器
props.put("value.serializer", "com.yourpackage.CustomMessageSerializer"); // 替换成你的类包路径

KafkaProducer<String, CustomMessage> producer = new KafkaProducer<>(props);

// 发送时直接传入CustomMessage对象
CustomMessage message = new CustomMessage(randomNumber, timestamp);
ProducerRecord<String, CustomMessage> record = new ProducerRecord<>("java-topic", Integer.toString(1), message);
producer.send(record);

额外提醒:区分Kafka元数据时间戳和业务时间戳

你可能注意到ProducerRecord有一个带timestamp参数的构造方法:

ProducerRecord(String topic, Integer partition, Long timestamp, K key, V value)

这个timestamp是Kafka存储的元数据时间戳,默认是Broker接收消息的时间,或者Producer发送消息的时间。如果你需要的是业务层面的时间戳(比如你生成随机数的时间),一定要把它放到消息体里,而不是用这个元数据字段——两者的用途完全不同。

内容的提问来源于stack exchange,提问作者Phil Baines

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 07:19:55