Ilya Shishkov created IGNITE-16176:
--------------------------------------
Summary: Make consumer request timeout configurable in
KafkaToIgniteCdcStreamerApplier
Key: IGNITE-16176
URL: https://issues.apache.org/jira/browse/IGNITE-16176
Project: Ignite
Issue Type: Improvement
Reporter: Ilya Shishkov
Now KafkaToIgniteCdcStreamerApplier performs requests with a hard-coded timeout:
{code:title=KafkaToIgniteCdcStreamerApplier#poll}
private void poll(KafkaConsumer<Integer, byte[]> cnsmr) throws
IgniteCheckedException {
ConsumerRecords<Integer, byte[]> recs =
cnsmr.poll(Duration.ofSeconds(DFLT_REQ_TIMEOUT));
if (log.isDebugEnabled()) {
log.debug(
"Polled from consumer [assignments=" + cnsmr.assignment() +
",rcvdEvts=" + rcvdEvts.addAndGet(recs.count()) + ']'
);
}
apply(F.iterator(recs, this::deserialize, true, rec ->
F.isEmpty(caches) || caches.contains(rec.key())));
cnsmr.commitSync(Duration.ofSeconds(DFLT_REQ_TIMEOUT));
}
{code}
We should have configurable timeout for requests to the Kafka.
--
This message was sent by Atlassian Jira
(v8.20.1#820001)