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

Fluentd-Kafka频繁连接关闭排查:AWS MSK高CPU异常

High Connection Creation/Close Rates with Fluentd Kafka2 Plugin on AWS MSK

We're running into an issue where our ~1000 EC2 instances using the Fluentd Kafka2 plugin are generating extremely high aws.kafka.connection_creation_rate and aws.kafka.connection_close_rate metrics (around 500 per minute) on our AWS MSK cluster (Kafka 3.3.1, 3 m4.xlarge nodes, multi-AZ, replication factor 3). Connections should ideally stay persistent, but in contrast, when simulating 1000 clients with kafka-producer-perf-test.sh (with message/byte throughput matching or exceeding production), the connection_close_rate drops to just 25 per minute, and cluster CPU usage is only half of production levels. No Fluentd connection-related warnings are logged.

Here's our Fluentd Kafka2 plugin configuration:

<label @default_kafka >
  <match *.**>
    @type copy
  <store>
    @type kafka2

    # list of seed brokers
    brokers {{ fluentd_env_vars[env]["xyz"] }}
    use_event_time true

    ssl_ca_cert "/path/CA_cert.pem"
    ssl_client_cert "/path/client_cert.pem"
    ssl_client_cert_key "/path/client_key.pem"
    
  # buffer settings
    <buffer topic>
      @type file
      path /var/log/td-agent/buffer/kafka
      flush_interval 1s
    </buffer>
      # data type settings
      <format>
        @type json
      </format>

      # topic settings
      topic_key topic
      default_topic log-messages

      # producer settings
      required_acks -1
      compression_codec gzip
      </store>
  </match>
  </label>

Root Cause Analysis & Fixes

Let's walk through the most likely reasons for this connection churn and how to address them:

1. Overly Aggressive Buffer Flush Interval

Your flush_interval is set to 1 second. This means each Fluentd instance is trying to send data to Kafka every single second. For 1000 instances, this creates a massive number of small flush operations—and if the plugin isn't reusing connections properly, each flush could trigger a new connection setup/teardown cycle. The producer perf test avoids this because it's optimized for sustained throughput with longer-lived connections.

Fix:
Increase the flush interval to a more reasonable value (e.g., 10–30 seconds) and tune buffer settings to batch more messages per flush. This reduces the number of connection attempts and encourages reuse:

<buffer topic>
  @type file
  path /var/log/td-agent/buffer/kafka
  flush_interval 10s  # Longer interval to batch more messages
  chunk_limit_size 100MB  # Larger chunks per flush
  queue_limit_length 32  # Handle more buffered chunks if needed
  flush_thread_count 4  # Parallelize flushes to avoid backpressure
</buffer>

2. Missing Connection Pooling Settings

The Kafka2 plugin supports connection pooling via max_connection and connection_max_idle_time, but these are missing from your config. Without pooling, the plugin might be creating a new connection for each flush instead of reusing existing idle connections.

Fix:
Add these connection pooling parameters to your Kafka2 store configuration to keep connections alive longer:

@type kafka2
# ... existing SSL/broker settings ...
max_connection 20  # Adjust based on your EC2 instance's resource capacity
connection_max_idle_time 300s  # Keep idle connections alive for 5 minutes

3. SSL Session Reuse Isn't Enabled

While you're using SSL for authentication, the plugin might not be reusing SSL sessions. Full SSL handshakes are expensive, and if each connection requires a new handshake, the plugin might be more likely to tear down connections instead of keeping them open.

Fix:
Enable SSL session reuse and extend SSL timeout:

ssl_session_reuse true
ssl_timeout 300s

4. Tune Kafka Producer Underlying Settings

The Kafka producer client used by Fluentd has default settings that might be causing premature connection closure. For example, connections_max_idle_ms (default 9 minutes in Kafka 3.3.1) could be too low when combined with a short flush interval. We can override these via producer_extra_params.

Fix:
Add these producer-specific settings to keep connections alive longer and reduce metadata refreshes:

producer_extra_params {
  connections_max_idle_ms: 1800000  # 30 minutes instead of 9
  metadata_max_age_ms: 300000        # Refresh metadata every 5 minutes instead of 30
}

5. Outdated Fluentd Kafka Plugin

Older versions of fluent-plugin-kafka have known bugs related to connection reuse. If you're running an outdated version, upgrading could resolve the issue.

Fix:
Update the plugin using td-agent's gem manager:

td-agent-gem update fluent-plugin-kafka

Post-Fix Verification

After applying these changes:

  • Monitor CloudWatch metrics for aws.kafka.connection_creation_rate and aws.kafka.connection_close_rate—they should drop to levels closer to what you saw with the producer perf test.
  • Check Fluentd logs for any new warnings/errors related to buffer flushes or connections (to ensure you aren't introducing backpressure).
  • Compare MSK cluster CPU usage to your perf test baseline—you should see a significant reduction.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 16:45:29