【发布时间】:2022-01-12 18:09:31
【问题描述】:
我正在尝试向 kafka 写入一条大消息(大约 15mb),但它没有被写入,程序结束,好像一切正常,但主题内没有消息。
确实会写一些小消息。
代码如下:
import org.apache.kafka.clients.producer.*;
import org.apache.kafka.common.serialization.StringSerializer;
import java.io.IOException;
import java.nio.file.Files;
import java.nio.file.Paths;
import java.util.Properties;
import java.util.concurrent.ExecutionException;
public class Main {
private final static String TOPIC = "rpdc_21596_in2";
private final static String BOOTSTRAP_SERVERS = "host:port";
private static KafkaProducer<String, String> createProducer() {
Properties props = new Properties();
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, BOOTSTRAP_SERVERS);
props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG,
StringSerializer.class.getName());
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG,
StringSerializer.class.getName());
props.put(ProducerConfig.MAX_REQUEST_SIZE_CONFIG, "20971520");
props.put("test.whatever", "fdsfdsf");
return new KafkaProducer<>(props);
}
public static void main(String[] args) throws ExecutionException, InterruptedException, IOException {
ProducerRecord<String, String> record =
new ProducerRecord<String, String>(TOPIC,
0,
123L,
"fdsfdsdsdssss",
new String(Files.readAllBytes(Paths.get("/Users/user/Desktop/value1.json")))
);
KafkaProducer<String, String> producer = createProducer();
RecordMetadata recordMetadata = producer.send(record).get();
producer.flush();
producer.close();
System.out.println(recordMetadata);
}
}
主题已配置为接受大消息,我已经能够用 python 写入它。这是python代码:
from kafka import KafkaProducer
from kafka.errors import KafkaError
producer = KafkaProducer(bootstrap_servers=['host:port'], max_request_size=20971520, request_timeout_ms=100000)
with open('/Users/user/Desktop/value1.json', 'rb') as f:
lines = f.read()
print(type(lines))
# produce keyed messages to enable hashed partitioning
future = producer.send('rpdc_21596_in2', key=b'foo', value=lines)
# Block for 'synchronous' sends
try:
record_metadata = future.get(timeout=50)
except KafkaError:
# Decide what to do if produce request failed...
pass
# Successful result returns assigned partition and offset
print (record_metadata.topic)
print (record_metadata.partition)
print (record_metadata.offset)
producer.flush()
但是那个java版本不行。
【问题讨论】:
-
recordMetadata 说什么,是否为空?
-
@ThomasRaffelsieper 不,它不为空。
-
尝试查看响应: producer.send(record, new Callback() { @Override public void onCompletion(RecordMetadata metadata, Exception exception) { // 如果 Exception 为 null,则记录发送成功} });
-
@ThomasRaffelsieper 是的,它正在工作。问题不在于代码。它是偏移资源管理器。它没有显示我的数据,但在主题的属性中,我只是注意到消息的数量确实增加了......谢谢伙计。
-
@OneCricketeer 我们团队的每个人都明白这一点,不知道他们为什么这样做,我认为这是一个临时解决方案
标签: java apache-kafka