DUE TO SPAM, SIGN-UP IS DISABLED. Goto Selfserve wiki signup and request an account.
...
| Code Block | ||
|---|---|---|
| ||
// Example 1: RecordTooLargeException use case
public class Example1 {
public static void main(String[] args) {
Properties producerProps = new Properties();
..... // omitted for brevity
producerProps.put("custom.exception.handler", Example1.ContinueTransaction.class.getName());
KafkaProducer<String, String> producer = new KafkaProducer(producerProps);
Properties consumerProps = new Properties();
..... // omitted for brevity
KafkaConsumer<String, String> consumer = new KafkaConsumer(consumerProps);
StringBuilder largeMessageStringBuffer = new StringBuilder();
for (int i = 0; i < 1000000; i++) {
largeMessageStringBuffer.append("0123456789");
}
String largeMessge = largeMessageStringBuffer.toString();
producer.initTransactions();
consumer.subscribe(Collections.singleton("input-topic"));
while(true) {
try {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofSeconds(60));
if (records.count() > 0) {
producer.beginTransaction();
Map<TopicPartition, OffsetAndMetadata> offsets = new HashMap();
for (ConsumerRecord<String, String> record : records) {
ProducerRecord<String, String> largeRecord = new ProducerRecord("output-topic", largeMessge);
Future<RecordMetadata> send = producer.send(largeRecord);
offsets.put(new TopicPartition(record.topic(), record.partition()), new OffsetAndMetadata(record.offset() + 1));
producer.send(new ProducerRecord("output-topic", "normalMessage");
offsets.put(new TopicPartition(record.topic(), record.partition()), new OffsetAndMetadata(record.offset() + 1));
}
producer.sendOffsetsToTransaction(offsets, consumer.groupMetadata());
producer.commitTransaction();
}
} catch (Exception e) {
producer.abortTransaction();
throw new RuntimeException(e);
}
}
}
public static class ContinueTransaction implements ProducerExceptionHandler {
@Override
public ProducerExceptionHandler.Response handle(ProducerRecord<byte[], byte[]> producerRecord, Exception e) {
return ProducerExceptionHandler.Response.SWALLOW;
}
@Override
public void configure(Map<String, ?> map) {
}
}
} |
...