使用Java Kafka客户端生产完所有消息后如何向主题发送元数据?
实现思路
要保证最后那条计数元数据消息是本次生产流程的最后一条,核心要满足两个条件:一是所有学生消息必须全部得到Kafka的成功ack之后,才能发送元数据消息;二是同批次消息要路由到同一个分区,因为Kafka只保证单个分区内的消息有序,跨分区没有顺序承诺。
具体实现步骤
- 首先配置生产者参数,从配置层面保证消息顺序和发送可靠性
Properties producerProps = new Properties(); producerProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka集群连接地址"); producerProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); // 如果学生对象用自定义序列化器就替换成对应类,用JSON序列化的话直接用StringSerializer即可 producerProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); // 开启幂等性,避免重试导致消息重复、乱序 producerProps.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true); producerProps.put(ProducerConfig.ACKS_CONFIG, "all"); // 幂等模式下该值设置不超过5即可保证单分区消息顺序 producerProps.put(ProducerConfig.MAX_IN_FLIGHT_REQUESTS_PER_CONNECTION, 5); KafkaProducer<String, String> producer = new KafkaProducer<>(producerProps);
- 发送时固定同批次消息的路由规则,保证所有学生消息和最后的元数据消息落在同一个分区
可以直接指定固定分区号,也可以给这批消息传相同的消息key,让默认分区器把它们路由到同一个分区,示例用固定分区的方式,逻辑更可控。 - 先批量提交所有学生消息的发送请求,阻塞等待所有消息发送成功后,再发送最后一条元数据消息
String topic = "你的目标主题名称"; // 固定发往0号分区,可根据业务调整 Integer fixedPartition = 0; List<Future<RecordMetadata>> sendFutures = new ArrayList<>(); // 遍历发送所有学生数据 for (Student stu : allStudent) { // 如果用JSON序列化,先把Student对象转成JSON字符串再传入 ProducerRecord<String, String> stuRecord = new ProducerRecord<>(topic, fixedPartition, "stu_data", JSON.toJSONString(stu)); sendFutures.add(producer.send(stuRecord)); } // 阻塞等待每一条学生消息发送成功,只要有一条失败就直接终止流程,避免发送错误的元数据 for (Future<RecordMetadata> future : sendFutures) { try { future.get(); } catch (InterruptedException | ExecutionException e) { producer.close(); throw new RuntimeException("学生数据发送异常,流程终止", e); } } // 所有学生消息确认成功后,发送最后一条计数消息 String finalMetaMsg = allStudent.size() + " records successfully produced to topic"; ProducerRecord<String, String> metaRecord = new ProducerRecord<>(topic, fixedPartition, "produce_meta", finalMetaMsg); // 同步等待元数据消息发送完成 producer.send(metaRecord).get(); // 所有流程结束再关闭生产者 producer.close();
如果业务要求跨分区场景下,消费者也绝对不会先读到元数据消息、后读到学生数据,可以开启生产者事务,把所有学生消息和元数据消息包裹在同一个事务中提交,消费者端设置
isolation.level = read_committed即可,只有事务提交成功后消费者才能读到这批消息,不会出现顺序错乱的问题。
内容的提问来源于stack exchange,提问作者Ajit Barik
相关产品推荐
相关产品推荐

