KafkaReceiver not running parallel for each partition?

Viewed 47

I'm trying to do read the messages concurrently with partition-based Ordering. I referred to the below example and I added delay (10sec) to mimic the workload. but I have noticed even though it runs multiple threads, it's not running concurrently. It only runs one message by message.

My topic consists of 20 partitions so I'm expecting to run at least 20 threads parallel and wait 10 sec (delay I put) and then run again and again.

enter link description here

    private fun onPayment(it: ReceiverRecord<String, Any>): Mono<ReceiverRecord<String, Any>> {
    println(
        "SMF message [${Thread.currentThread().name}], Offset -  ${
        it.receiverOffset().offset()
        }  , Partition ${
        it.receiverOffset().topicPartition()
        },  Value - ${it.value()}  ${SimpleDateFormat("hh:mm:ss").format(Date())}"
    )
    return Mono.just(it).delayElement(Duration.ofMillis(10000))
}

fun doSubscribe(consumer: (ReceiverRecord<String, Any>) -> Mono<ReceiverRecord<String, Any>>) {
    val reOption = KafkaConfig().kafkaConsumerConfig().subscription(Collections.singleton("payTopic"))
        .commitInterval(Duration.ZERO)
        .commitBatchSize(0)
    val receiver = KafkaReceiver.create(reOption)

    Flux.defer(receiver::receive).groupBy { m ->
        m.receiverOffset().topicPartition()
    }.flatMap { grpFlux ->
        grpFlux.publishOn(schd).concatMap { m ->
            consumer(m).map {
                m.receiverOffset().commit()
            }
        }
    }.subscribe()
}

And its prints sequentially one by one and consume only from one partition.

enter image description here

0 Answers
Related