【问题标题】:Losing some events in Kafka Joins在 Kafka Joins 中丢失一些事件
【发布时间】:2021-09-03 12:42:50
【问题描述】:

我的 Spring-Cloud 应用程序中有 3 个不同的流。 每天有数十万条记录,然而,每天大约有 3 或 4 条记录消失。 我看到它们出现在日志中,但是,并没有完成所有的连接。

代码:

  @StreamListener
fun processOrderEvent(
    @Input(StreamBindings.BUSINESS_CONDITION_ORDER_CHARGED_CHARGE_IN)
    chargeEvent: KStream<String, ChargeEvent>,
    @Input(StreamBindings.BUSINESS_CONDITION_ORDER_CHARGED_IN)
    orderEvent: KStream<String, OrderChargedEvent>,
    @Input(StreamBindings.BUSINESS_CONDITION_ORDER_CREATED_CHARGE_IN)
    orderCreatedEvent: KStream<String, OrderCreatedEvent>
) {

        val tracing = Tracing.newBuilder().build()
        val kafkaStreamsTracing = KafkaStreamsTracing.create(tracing)

        val chargeKeyValue = chargeEvent
            .filter { _, event -> shouldProcessBusinessConditions(event) && (event.flow == "___________" || event.flow == "____")}
            .transform(
                kafkaStreamsTracing.map<String, ChargeEvent, ByteArray, ByteArray>("processOrderEvent_ChargeEvent") { _, value ->
                    var traceId = tracing.tracer().currentSpan().context().traceIdString()
                    val keyValue = KeyValue(value.id.toString().toByteArray(), objectMapper.writeValueAsString(
                        Charge(
                            value.id.toString(),
                            value.amount?.value,
                            value.status?.name,
                            LocalDateTime.now().toString(),
                            value.paymentMethod?.type?.name,
                            value.creditor?.customerId,
                            value.channel?.name,
                            value.amount?.currency?.name,
                            value.paymentMethod?.installments,
                            value.card?.brand,
                            buildSellerEmail(value),
                            value.createdAt,
                            value.amount?.summary?.total,
                            value.amount?.summary?.paid,
                            value.amount?.summary?.refunded,
                            value.connect?.id,
                            value.connect?.name,
                            value.flow
                        )).toByteArray())

                    log.info("m=processOrderEvent traceId=$traceId chargeId=${value.id} step=chargeKeyValue")
                    keyValue
                }
            )

        val orderKeyValue = orderEvent
            .transform(
                kafkaStreamsTracing.map<String, OrderChargedEvent, ByteArray, ByteArray>("processOrderEvent_OrderChargedEvent") { _, value ->
                    var traceId = tracing.tracer().currentSpan().context().traceIdString()
                    log.info("m=processOrderEvent traceId=$traceId chargeId=${value.chargeId} orderId=${value.orderId} step=orderKeyValue")
                    KeyValue(value.chargeId.toByteArray(), objectMapper.writeValueAsString(Order(value.orderId, value.chargeId)).toByteArray())
                }
            )

        val orderCreatedKeyValue = orderCreatedEvent
            .transform(
                kafkaStreamsTracing.map<String, OrderCreatedEvent, ByteArray, ByteArray>("processOrderEvent_OrderCreatedEvent") { _, value ->
                    var traceId = tracing.tracer().currentSpan().context().traceIdString()
                    log.info("m=processOrderEvent traceId=$traceId orderId=${value.orderId} step=before_orderCreatedKeyValue")
                    val originalValue = buildOrderOriginalValue(value)
                    val keyValue = KeyValue(value.orderId.toByteArray(), objectMapper.writeValueAsString(OrderCreated(value.orderId, originalValue)).toByteArray())
                    log.info("m=processOrderEvent traceId=$traceId orderId=${value.orderId} originalValue=${originalValue} step=orderCreatedKeyValue")
                    keyValue
                }
            )

        chargeKeyValue.join(orderKeyValue, OrderChargeValueJoiner(), JoinWindows.of(Duration.ofHours(5)))
            .transform(
                kafkaStreamsTracing.map<ByteArray, ByteArray, ByteArray, ByteArray>("processOrderEvent_OrderChargedEvent_ChargeEvent") { _, value ->
                    var traceId = tracing.tracer().currentSpan().context().traceIdString()
                    val orderWithChargeJson = objectMapper.readValue(value, OrderWithCharge::class.java)
                    val keyValue = KeyValue(orderWithChargeJson.order.orderId!!.toByteArray(), value)
                    log.info("m=processOrderEvent traceId=$traceId chargeId=${orderWithChargeJson.charge.chargeId} orderId=${orderWithChargeJson.order.orderId} step=orderKeyValueJoin")
                    keyValue
                }
            )
            .join(orderCreatedKeyValue, OrderCreatedValueJoiner(), JoinWindows.of(Duration.ofHours(5)))
            .transform(
                kafkaStreamsTracing.map<ByteArray, ByteArray, ByteArray, ByteArray>("processOrderEvent_OrderChargedEvent_ChargeEvent_OrderCreatedEvent") { key, value ->
                    var traceId = tracing.tracer().currentSpan().context().traceIdString()
                    val event = objectMapper.readValue(value, OrderWithChargeAndOrderCreated::class.java)
                    log.info("m=processOrderEvent traceId=$traceId chargeId=${event.charge.chargeId} orderId=${event.order.orderId} step=orderCreatedKeyValueJoin")
                    KeyValue(key, objectMapper.writeValueAsString(OrderWithChargeAndOrderCreatedTraceId(traceId, event)).toByteArray())
                }
            ).process(ProcessorSupplier { businessConditionCkoutEventProcessor })
}

少数无效之一的日志:

21 年 9 月 1 日 上午 5:48:15.863
e2e64c0faa9e 05:48:15.863 [orders-chargeds-charges-v3-df80aa80-d3f6-48c7-a862-7a05c55d5d24-StreamThread-1] 信息 u.p.p.s.SomeEventStreams - m=processOrderEvent traceId=06352e7dfb290f6fchargeId=--c /em>-58ecdbf7ccf7 orderId=ORDE__-413D-4D5C-_-1873F0DB5ADE step=orderKeyValueJoin

21 年 9 月 1 日 上午 5:48:15.763
c232aa303f2f 05:48:15.763 [orders-chargeds-charges-v3-a1f58fb2-63d6-4195-bffd-6a9009b00707-StreamThread-1] 信息 u.p.p.s.SomeEventStreams - m=processOrderEvent traceId=06352e7dfb290f84f chargeId=-c /em>-58ecdbf7ccf7 orderId=ORDE__-413D-4D5C-_-1873F0DB5ADE step=orderKeyValue

21 年 9 月 1 日 上午 5:48:15.749
c232aa303f2f 05:48:15.749 [orders-chargeds-charges-v3-a1f58fb2-63d6-4195-bffd-6a9009b00707-StreamThread-1] 信息 u.p.p.s.SomeEventStreams - m=processOrderEvent traceId=7f8c7bb5bbdb422b4891=-c71 /em>-58ecdbf7ccf7 step=chargeKeyValue

21 年 8 月 31 日 下午 6:31:50.499
c232aa303f2f 18:31:50.499 [orders-chargeds-charges-v3-a1f58fb2-63d6-4195-bffd-6a9009b00707-StreamThread-1] 信息 u.p.p.s.SomeEventStreams-m=processOrderEvent traceId=95fc78a495f173F64 orderId=ORDE--1 originalValue=22000 step=orderCreatedKeyValue

21 年 8 月 31 日 下午 6:31:50.499
c232aa303f2f 18:31:50.499 [orders-chargeds-charges-v3-a1f58fb2-63d6-4195-bffd-6a9009b00707-StreamThread-1] INFO u.p.p.s.SomeEventStreams - m=processOrderEvent traceId=95fc78a495f03b64 orderId=ORDE_3_-413D-4D5C-_-1873F0DB5ADE step=before_orderCreatedKeyValue

数千个有效之一的日志:

21 年 9 月 2 日 下午 6:34:33.547
f84980d99867 18:34:33.547 [orders-chargeds-charges-v3-63c5bb8a-8026-4b05-937f-c41360f90201-StreamThread-1] 信息 u.p.p.s.SomeEventStreams - m=processOrderEvent traceId=3b5981d48ac7e3a-emId=- >-8576-3cb1173804d5 orderId=ORDE__-084F-4CCE-_-BA4DCAF1D50A step=orderCreatedKeyValueJoin

21 年 9 月 2 日 下午 6:34:33.446
e2e64c0faa9e 18:34:33.446 [orders-chargeds-charges-v3-df80aa80-d3f6-48c7-a862-7a05c55d5d24-StreamThread-1] INFO u.p.p.s.SomeEventStreams - m=processOrderEvent traceId=3b5981d48ac7e130 chargeId=3834ad80-243a-4ffe-8576-3cb1173804d5 orderId=ORDE__-084F-4CCE-_-BA4DCAF1D50A step=orderKeyValueJoin

21 年 9 月 2 日 下午 6:34:33.346
f84980d99867 18:34:33.346 [orders-chargeds-charges-v3-63c5bb8a-8026-4b05-937f-c41360f90201-StreamThread-1] 信息 u.p.p.s.SomeEventStreams - m=processOrderEvent traceId=6501784c0bf584D54CID=ORDE-BA originalValue=85400 step=orderCreatedKeyValue

21 年 9 月 2 日 下午 6:34:33.346
f84980d99867 18:34:33.346 [orders-chargeds-charges-v3-63c5bb8a-8026-4b05-937f-c41360f90201-StreamThread-1] 信息 u.p.p.s.SomeEventStreams - m=processOrderEvent traceId=6501784c0bf584D54CID=ORDE-BA step=before_orderCreatedKeyValue

21 年 9 月 2 日 下午 6:34:33.344
c232aa303f2f 18:34:33.344 [orders-chargeds-charges-v3-a1f58fb2-63d6-4195-bffd-6a9009b00707-StreamThread-1] 信息 u.p.p.s.SomeEventStreams - m=processOrderEvent traceId=3b5981d48ac7e130charge-Id=-2 /em>-3cb1173804d5 orderId=-084F-4CCE--BA4DCAF1D50A step=orderKeyValue

21 年 9 月 2 日 下午 6:34:33.326
c232aa303f2f 18:34:33.326 [orders-chargeds-charges-v3-a1f58fb2-63d6-4195-bffd-6a9009b00707-StreamThread-1] 信息 u.p.p.s.SomeEventStreams - m=processOrderEvent traceId=efd790c85a8f053c chargeId=-4 >-3cb1173804d5 step=chargeKeyValue

“OrderCreatedKeyValuejoin”日志是结尾,在大多数情况下都会显示,但在某些情况下,事件永远不会在结尾出现。

【问题讨论】:

    标签: spring-boot apache-kafka stream apache-kafka-streams spring-cloud-stream


    【解决方案1】:

    您需要检查如何使用 Kafka 实现 DLQ。

    要实现 DLQ,您需要创建单独的队列,并且无论主题失败,您都需要将其放入 DLQ。编写重试逻辑来读取它。

    也许这些文件对你有帮助 https://cloud.spring.io/spring-cloud-static/spring-cloud-stream-binder-kafka/2.1.0.RC1/multi/multi_kafka-dlq-processing.html

    Dead letter queue (DLQ) for Kafka with spring-kafka

    【讨论】:

    • 所有被调用的方法都有日志和try/catchs。如果按照您所说的发生任何错误,它会显示在日志中,不是吗?我认为joinwindows之间存在一些差距
    猜你喜欢
    • 1970-01-01
    • 2019-08-10
    • 2023-03-09
    • 2020-02-13
    • 2022-01-18
    • 1970-01-01
    • 1970-01-01
    • 2017-10-30
    • 2021-01-05
    相关资源
    最近更新 更多