Ash*_*iya 2 python java apache-kafka spring-kafka
在我们的应用程序中,生产者和消费者都启用了一次。
Producer是一个python组件。我们已经启用:
Consumer 是一个 Spring Boot 应用程序。我们已启用:
我们在 ConfluenceCloud 上有多分区 Kafka 主题(假设有 3 个分区)。
我们的应用程序设计如下:
问题:
我们注意到,有时同一条 Kafka 消息会在 Consumer 中被多次消费。我们通过使用以下消费者代码检测到了这一点。我们将之前消费的kafka消息Id(带偏移量)保存在Redis中,并与新消费的消息进行比较。
消费者代码:
@KafkaListener(topics = "${datalake.datasetevents.topic}", groupId = "${spring.kafka.consumer.group-id}")
public void listen(@Header(KafkaHeaders.RECEIVED_MESSAGE_KEY) String key,
@Header(KafkaHeaders.OFFSET) String offset,
@Payload InputEvent inputEvent, Acknowledgment acknowledgment) {
//KafkaHeaders.
Event event = new Event();
event.setCorrId(inputEvent.getCorrId());
event.setQn(inputEvent.getQn());
event.setCreatedTs(new Date());
event.setEventTs(inputEvent.getEventTs());
event.setMeta(inputEvent.getMeta() != null ? inputEvent.getMeta(): new HashMap<>());
event.setType(inputEvent.getType());
event.setUlid(key);
//detect message duplications
try {
String eventRedisKey = "tg_e_d_" + key.toLowerCase();
String redisVal = offset;
String tmp = redisTemplateString.opsForValue().get(eventRedisKey);
if (tmp != null) {
dlkLogging.error("kafka_event_dup", "Event consumed more than once ulid:" + event.getUlid()+ " redis offset: "+tmp+ " event offset:"+offset);
redisTemplateString.delete(eventRedisKey);
}
redisTemplateString.opsForValue().set(eventRedisKey, redisVal, 30, TimeUnit.SECONDS);
} catch (Exception e) {
dlkLogging.error("kafka_consumer_redis","Redis error at kafka consumere ", e);
}
//process the message and ack
try {
eventService.saveEvent(persistEvent, event);
ack.acknowledge();
} catch (Exception ee) {
//Refer : /sf/ask/4368928931/
ack.nack(1);
dlkLogging.error("event_sink_error","error sinking kafka event.Will retry", ee);
}
}
Run Code Online (Sandbox Code Playgroud)
行为:我们注意到“kafka_event_dup”每天发送多次。
错误消息:事件消耗多次 ulid:01G77G8KNTSM2Q01SB1MK60BTH redis 偏移量:659238 事件偏移量:659238
问题: 尽管我们在生产者和消费者中都配置了一次完全相同的消息,为什么消费者会读取相同的消息?
更新:在阅读了几篇SO帖子后,即使配置了exactly-once,我们似乎仍然需要在消费者端实现重复数据删除逻辑?
附加信息:
消费者配置:
public DefaultKafkaConsumerFactory kafkaDatasetEventConsumerFactory(KafkaProperties properties) {
Map<String, Object> props = properties.buildConsumerProperties();
props.put(ENABLE_AUTO_COMMIT_CONFIG, false);
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, ErrorHandlingDeserializer.class);
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, ErrorHandlingDeserializer.class);
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, ErrorHandlingDeserializer.class);
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, ErrorHandlingDeserializer.class);
props.put(ConsumerConfig.ISOLATION_LEVEL_CONFIG, "read_committed");
props.put(ErrorHandlingDeserializer.KEY_DESERIALIZER_CLASS, StringDeserializer.class);
props.put(ErrorHandlingDeserializer.VALUE_DESERIALIZER_CLASS, CustomJsonDeserializer.class.getName());
props.put(JsonDeserializer.VALUE_DEFAULT_TYPE, "com.fr.det.datalake.eventdriven.model.kafka.InputEvent");
return new DefaultKafkaConsumerFactory(props);
}
Run Code Online (Sandbox Code Playgroud)
生产者代码(python):
def __get_producer(self):
conf = {
'bootstrap.servers': self.server,
'enable.idempotence': True,
'acks': 'all',
'retry.backoff.ms': self.sleep_seconds * 100
}
if self.sasl_mechanism:
conf['sasl.mechanisms'] = self.sasl_mechanism
if self.security_protocol:
conf['security.protocol'] = self.security_protocol
if self.sasl_username:
conf['sasl.username'] = self.sasl_username
if self.sasl_username:
conf['sasl.password'] = self.sasl_password
if self.transaction_prefix:
conf['transactional.id'] = self.__get_transaction_id()
producer = Producer(conf)
return producer
@_retry_on_error
def send_messages(self, messages, *args, **kwargs):
ts = time.time()
producer = kwargs.get('producer', None)
if producer is not None:
for message in messages:
key = message.get('key', str(ulid.from_timestamp(ts)))
value = message.get('value', None)
topic = message.get('topic', self.topic)
producer.produce(topic=topic,
value=value,
key=key,
on_delivery=self.acked)
producer.commit_transaction(30)
def _retry_on_error(func, *args, **kwargs):
def inner(self, messages, *args, **kwargs):
attempts = 0
while True:
attempts += 1
sleep_time = attempts * self.sleep_seconds
try:
producer = self.__get_producer()
self.logger.info(f"Producer: {producer}, Attempt: {attempts}")
producer.init_transactions(30)
producer.begin_transaction()
res = func(self, messages, *args, producer=producer, **kwargs)
return res
except KafkaException as e:
if attempts <= self.retry_count:
if e.args[0].txn_requires_abort():
producer.abort_transaction(30)
time.sleep(sleep_time)
continue
self.logger.error(str(e), exc_info=True, extra=extra)
break
return inner
Run Code Online (Sandbox Code Playgroud)
Kafka 的精确一次本质上是 Kafka-Streams 的一项功能,尽管它也可以与常规消费者和生产者一起使用。
精确一次只能在应用程序仅与 Kafka 交互的环境中实现:没有 XA 或其他类型的跨技术分布式事务可以使 Kafka 消费者与其他存储(如 Redis)进行一次精确交互方式。
在分布式世界中,我们必须承认这是不可取的,因为它会引入锁定、争用以及负载下性能呈指数级下降。如果我们不需要处于分布式世界,那么我们就不需要 Kafka,很多事情都会变得更容易。
Kafka 中的事务旨在在仅与 Kafka 交互的一个应用程序中使用,它可以让您保证应用程序将 1) 从某些主题分区读取,2) 将一些结果写入其他一些主题分区,3) 提交读取与 1 相关的偏移量,或者不执行这些操作。如果多个应用程序以这种方式背对背放置并通过 Kafka 进行交互,那么如果您非常小心,您就可以实现一次。如果您的消费者需要 4) 与 Redis 交互 5) 与其他存储交互或在某处执行一些副作用(例如发送电子邮件等),那么通常无法执行步骤 1,2,3,4, 5 以原子方式作为分布式应用程序的一部分。你可以用其他存储技术来实现这种事情(是的,Kafka本质上就是一种存储),但它们不能分布式,你的应用程序也不能。这本质上就是 CAP 定理告诉我们的。
这也是为什么恰好一次本质上是 Kafka 流的东西:Kafka Stream 只是 Kafka 消费者/生产者客户端的智能包装器,以免您构建仅与 Kafka 交互的应用程序。
您还可以使用其他数据处理框架(例如 Spark Streaming 或 Flink)实现一次性流处理。
在实践中,不关心事务而只在消费者中进行重复数据删除通常要简单得多。您可以保证消费者组中最多有一个消费者在任何时间点连接到每个分区,因此重复的情况总是会发生在您的应用程序的同一实例中(直到它重新扩展),并且取决于您的配置,复制通常应该只发生在一个 Kafka 消费者缓冲区内,因此您不需要在消费者中存储太多状态来进行重复数据删除。如果您使用某种只能增加的事件 ID(顺便说一句,这本质上就是 Kafka 偏移量,这并非巧合),那么您只需要在应用程序的每个实例的状态中保留最大事件-您已成功处理的每个分区的 id。
| 归档时间: |
|
| 查看次数: |
1862 次 |
| 最近记录: |