From ed77484f72a46435dd6b0c67ac91075b23825fcc Mon Sep 17 00:00:00 2001 From: Christian Tzolov Date: Wed, 9 Jun 2021 15:19:19 +0200 Subject: [PATCH] Try to mitigate cdc overflow error --- .../cloud/fn/supplier/cdc/CdcSupplierConfiguration.java | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/functions/supplier/cdc-debezium-supplier/src/main/java/org/springframework/cloud/fn/supplier/cdc/CdcSupplierConfiguration.java b/functions/supplier/cdc-debezium-supplier/src/main/java/org/springframework/cloud/fn/supplier/cdc/CdcSupplierConfiguration.java index ba546f23..5be4b899 100644 --- a/functions/supplier/cdc-debezium-supplier/src/main/java/org/springframework/cloud/fn/supplier/cdc/CdcSupplierConfiguration.java +++ b/functions/supplier/cdc-debezium-supplier/src/main/java/org/springframework/cloud/fn/supplier/cdc/CdcSupplierConfiguration.java @@ -31,6 +31,7 @@ import org.apache.kafka.connect.source.SourceRecord; import reactor.core.publisher.EmitterProcessor; import reactor.core.publisher.Flux; import reactor.core.publisher.FluxSink; +import reactor.core.publisher.Sinks; import org.springframework.beans.factory.BeanClassLoaderAware; import org.springframework.boot.context.properties.EnableConfigurationProperties; @@ -114,7 +115,7 @@ public class CdcSupplierConfiguration implements BeanClassLoaderAware { Function recordFlattening, ObjectMapper mapper, CdcSupplierProperties cdcStreamingEngineProperties) { - FluxSink> sink = emitterProcessor.sink(); + FluxSink> sink = emitterProcessor.sink(FluxSink.OverflowStrategy.BUFFER); Consumer messageConsumer = sourceRecord -> { // When cdc.flattening.deleteHandlingMode=none and cdc.flattening.dropTombstones=false