diff --git a/plugins/event-bus/kafka/src/main/java/org/apache/cloudstack/mom/kafka/KafkaEventBus.java b/plugins/event-bus/kafka/src/main/java/org/apache/cloudstack/mom/kafka/KafkaEventBus.java index f2589d2d7d09..daadb5a523d6 100644 --- a/plugins/event-bus/kafka/src/main/java/org/apache/cloudstack/mom/kafka/KafkaEventBus.java +++ b/plugins/event-bus/kafka/src/main/java/org/apache/cloudstack/mom/kafka/KafkaEventBus.java @@ -46,6 +46,7 @@ public class KafkaEventBus extends ManagerBase implements EventBus { public static final String DEFAULT_TOPIC = "cloudstack"; public static final String DEFAULT_SERIALIZER = "org.apache.kafka.common.serialization.StringSerializer"; + public static final String DEFAULT_MAX_BLOCK_MS = "2500"; private String _topic = null; private Producer _producer; @@ -70,6 +71,10 @@ public boolean configure(String name, Map params) throws Configu if (!props.containsKey("value.serializer")) { props.put("value.serializer", DEFAULT_SERIALIZER); } + + if (!props.containsKey("max.block.ms")) { + props.put("max.block.ms", DEFAULT_MAX_BLOCK_MS); + } } catch (Exception e) { throw new ConfigurationException("Could not read kafka properties"); }