You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在Kafka节点宕机时停止数据推送,避免Java Web服务器内存溢出

Solution for Stopping Kafka Message Production When Broker is Unavailable

Hey there, let's fix that memory crash issue you're hitting when Kafka goes down! The root problem here is that the default Kafka producer will keep buffering unsent messages in memory and retrying indefinitely—since your web server is churning out HTTP requests nonstop, that buffer just keeps growing until it takes down your server. Here's how to tackle this step by step:

1. Tune Producer Configs to Cap Memory Usage

First, adjust your producer properties to prevent it from hoarding too many messages in memory. These settings will limit buffer size, retry attempts, and how long the producer waits before failing:

static Producer<String, String> producer;

void initProducer() {
    Properties properties = new Properties();
    // Basic Kafka connection settings
    properties.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "your-kafka-node-ip:9092");
    properties.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
    properties.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());

    // Critical memory & retry limits
    properties.put(ProducerConfig.BUFFER_MEMORY_CONFIG, 67108864); // 64MB total buffer cap
    properties.put(ProducerConfig.MAX_BLOCK_MS_CONFIG, 30000); // Block for max 30s if buffer is full, then throw
    properties.put(ProducerConfig.RETRIES_CONFIG, 3); // Only retry 3 times
    properties.put(ProducerConfig.RETRY_BACKOFF_MS_CONFIG, 1000); // Wait 1s between retries
    properties.put(ProducerConfig.ACKS_CONFIG, "1"); // Balance durability and speed

    producer = new KafkaProducer<>(properties);
}

2. Add a Kafka Health Check

Before sending any messages, verify if Kafka is reachable. Use AdminClient to run a quick cluster health check:

private boolean isKafkaReachable() {
    try (AdminClient adminClient = AdminClient.create(properties)) {
        // Try to fetch cluster ID with a short timeout
        adminClient.describeCluster().clusterId().get(5, TimeUnit.SECONDS);
        return true;
    } catch (Exception e) {
        log.error("Kafka cluster is unavailable", e);
        return false;
    }
}

Then update your send logic to skip sending when Kafka is down:

public void processAndSend(String topic, String requestData) {
    // Process your HTTP request data here...
    String processedData = processRequestData(requestData);

    if (!isKafkaReachable()) {
        log.warn("Kafka is down—skipping message send for request data: {}", processedData);
        // Optional: Store the message in a persistent queue (like Redis/DB) to replay later
        queueForLaterProcessing(processedData);
        return;
    }

    ProducerRecord<String, String> record = new ProducerRecord<>(topic, processedData);
    producer.send(record, (metadata, exception) -> {
        if (exception != null) {
            log.error("Failed to send message to Kafka", exception);
            // Trigger a temporary pause or queue the message
            queueForLaterProcessing(processedData);
        }
    });
}

3. Use a Circuit Breaker to Stop Sending on Failures

For more robust failure handling, implement a circuit breaker (like Resilience4j) to automatically pause sending when consecutive failures happen:

// Initialize circuit breaker
CircuitBreakerConfig cbConfig = CircuitBreakerConfig.custom()
    .failureRateThreshold(50) // Open circuit if 50% of sends fail
    .waitDurationInOpenState(Duration.ofMinutes(1)) // Wait 1min before retrying
    .build();
CircuitBreaker kafkaCircuitBreaker = CircuitBreaker.of("kafka-producer", cbConfig);

public void sendWithCircuitBreaker(String topic, String data) {
    Try.run(() -> kafkaCircuitBreaker.executeSupplier(() -> {
        producer.send(new ProducerRecord<>(topic, data)).get();
        return null;
    }))
    .onFailure(e -> {
        log.error("Kafka send failed—circuit breaker triggered", e);
        queueForLaterProcessing(data);
    });
}

Key Additional Tips

  • Don't Drop Messages Blindly: When Kafka is down, store unsent messages in a persistent store (database, Redis) so you can replay them once Kafka recovers.
  • Proper Producer Shutdown: Always call producer.close() when your web server shuts down to flush any remaining messages safely.
  • Avoid Blocking Web Threads: Keep message sending async so your HTTP request handlers don't get stuck waiting for Kafka.

内容的提问来源于stack exchange,提问作者Juvenik

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.21 07:50:44