DUE TO SPAM, SIGN-UP IS DISABLED. Goto Selfserve wiki signup and request an account.
...
| Code Block | ||||||||
|---|---|---|---|---|---|---|---|---|
| ||||||||
package org.apache.kafka.clients.producer;
import org.apache.kafka.common.Configurable;
import org.apache.kafka.common.errors.RecordTooLargeException;
import org.apache.kafka.common.errors.UnknownTopicOrPartitionException;
import java.io.Closeable;
/**
* Interface that specifies how an the RecordTooLargeException and/or UnknownTopicOrPartitionException should be handled.
* The accepted responses for RecordTooLargeException are FAIL and SWALLOW. Therefore, RETRY will be interpreted and executed as FAIL.
*/
public interface ProducerExceptionHandler extends Configurable, Closeable {
/**
* Determine whether to stop processing, or swallow the error by dropping the record.
*
* @param record The record that failed to produce
* @param exception The exception that occurred during production
*/
default NonRetryableResponse handle(final ProducerRecord record, final RecordTooLargeException exception) {
return NonRetryableResponse.FAIL;
}
/**
* Determine whether to stop processing, keep retrying internally, or swallow the error by dropping the record.
*
* @param record The record that failed to produce
* @param exception The exception that occurred during production
*/
default URetryableResponseRetryableResponse handle(final ProducerRecord record, final UnknownTopicOrPartitionException exception) {
return RetryableResponse.RETRY;
}
enum NonRetryableResponse {
/* stop processing: fail */
FAIL(0, "FAIL"),
/* drop the record and continue */
SWALLOW(1, "SWALLOW");
/**
* an english description of the api--this is for debugging and can change
*/
public final String name;
/**
* the permanent and immutable id of an API--this can't change ever
*/
public final int id;
RecordTooLargeExceptionResponse(final int id, final String name) {
this.id = id;
this.name = name;
}
}
enum retryableResponse {
/* stop processing: fail */
FAIL(0, "FAIL"),
/* continue: keep retrying */
RETRY(1, "RETRY"),
/* continue: swallow the error */
SWALLOW(2, "SWALLOW");
/**
* an english description of the api--this is for debugging and can change
*/
public final String name;
/**
* the permanent and immutable id of an API--this can't change ever
*/
public final int id;
UnknownTopicOrPartitionExceptionResponse(final int id, final String name) {
this.id = id;
this.name = name;
}
}
} |
...
| Code Block | ||||||
|---|---|---|---|---|---|---|
| ||||||
// Example 1: RecordTooLargeException use case
// In this example we except that the producer follows the custom handler and does not fail. It may make a batch of normal records and commit the transaction successfully.
public class Example1 {
public static void main(String[] args) {
Properties producerProps = new Properties();
producerProps.put("custom.exception.handler.class", Example1.MyProducerExceptionHandler.class.getName());
// ..... omitted for brevity
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 largeMessage = largeMessageStringBuffer.toString();
ProducerRecord<String, String> largeRecord = new ProducerRecord("output-topic", largeMessage);
ProducerRecord<String, String> normalRecord = new ProducerRecord("output-topic", "normalMessage");
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) {
Future<RecordMetadata> sendOutput;
if (someMethod(record)) {
sendOutput = producer.send(largeRecord);
} else {
sendOutput = producer.send(normalRecord);
}
if (!sendOutput instanceof FutureFailure) {
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 MyProducerExceptionHandler implements ProducerExceptionHandler {
@Override
public ProducerExceptionHandler.NonRetryableResponse handle(ProducerRecord producerRecord, RecordTooLargeExceptionException e) {
return ProducerExceptionHandler.NonRetryableResponse.SWALLOW;
}
@Override
public ProducerExceptionHandler.RetryableResponsevoid handleconfigure(ProducerRecord producerRecordMap<String, UnknownTopicOrPartitionException?> emap) {
}
return null;@Override
}
@Override
public void configure(Map<String, ?> map) {
}
@Override
public void close() throws IOException {
}
}
} |
...
| Code Block | ||||||
|---|---|---|---|---|---|---|
| ||||||
// Example 2: UnknownTopicOrPartitionException use case
// In this example we except that if the record does NOT belong to the "important-topic", the producer follows the custom handler and fails.
public class Example2 {
public static void main(String[] args) {
Properties producerProps = new Properties();
producerProps.put("custom.exception.handler.class", Example2.MyProducerExceptionHandler.class.getName());
// ..... omitted for brevity
KafkaProducer<String, String> producer = new KafkaProducer(producerProps);
producer.send(new ProducerRecord(someMethodToComputeTopic(), "someMessage"));
producer.flush();
producer.close();
}
public static class MyProducerExceptionHandler implements ProducerExceptionHandler {
@Override
public ProducerExceptionHandler.RetryableResponse handle(ProducerRecord producerRecord, UnknownTopicOrPartitionException e) {
if (producerRecord.topic().equals("important-topic") {
ProducerExceptionHandler.RetryableResponse.RETRY;
}
return ProducerExceptionHandler.RetryableResponse.FAIL;
}
@Override
public ProducerExceptionHandler.NonRetryableResponse handle(ProducerRecord producerRecord, RecordTooLargeException e) {
return null;
}
@Override
public void configure(Map<String, ?> map) {
}
@Override
public void close() throws IOException {
}
}
} |
...