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