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

Docker环境下Fluentd日志可输出到stdout但无法写入Elasticsearch

Troubleshooting Your Fluentd → Elasticsearch Log Pipeline

Let's work through your issue step by step—you've got two core problems here: connection failures between Fluentd and Elasticsearch, and log parsing/mapping errors that prevent data from landing in Kibana. Here's how to fix them:

1. Fix Elasticsearch Connectivity Issues

Your earlier connection errors stem from how Docker containers communicate. If you're using Docker Compose, containers in the same default network can reach each other via service names, not just host IPs. Using 172.18.0.1 or docker.for.mac.localhost is unreliable here—instead, reference your Elasticsearch service directly by name (e.g., elasticsearch if that's what's in your docker-compose.yml).

Also, add retry/timeout settings to handle Elasticsearch's slow startup (it needs time to initialize before accepting connections).

2. Fix Log Parsing & Mapping Errors

The mapper_parsing_exception happens because you're sending a stringified JSON object from Python, but Elasticsearch expects structured data or plain text. When Fluentd tries to index this string as a nested log object, Elasticsearch throws an error because the field mapping doesn't match.

Updated Fluentd Configuration

Replace your current config with this (adjust service names if needed for your setup):

# INJECTED VIA DOCKER COMPOSE
<source>
  @type forward
  port 24224
  bind 0.0.0.0 # Ensure Fluentd accepts connections from outside its container
</source>

# Optional: Uncomment only if you still send stringified JSON from Python
# <filter **>
#   @type parser
#   format json
#   key_name message # Fluent-logger-python puts log content in `message` by default
#   hash_value_field parsed
#   reserve_data true
# </filter>

<match **>
  @type copy
  <store>
    @type elasticsearch
    # Use your Elasticsearch service name from docker-compose.yml
    hosts elasticsearch:9200
    logstash_format true
    logstash_prefix chris.risley
    logstash_dateformat %Y%m%d
    include_tag_key true
    flush_interval 1s
    # Add resilience for slow ES startup
    retry_limit 10
    retry_wait 2s
    timeout 10s
  </store>
  <store>
    @type stdout
    format json # Makes it easier to inspect structured logs
  </store>
</match>

Updated Python Logger Code

Instead of sending stringified JSON, send native Python dictionaries—this lets Fluentd pass structured data directly to Elasticsearch without parsing errors:

from fluent import handler
import logging
import msgpack
from io import BytesIO

class ElasticLogger(object):
    def __init__(self, tag, host, port, base_level, config=False, path='src/logging.yaml'):
        self.tag = tag
        self.host = host
        self.port = port
        self.base_level = base_level
        self.logger = self.config(config, path)

    def build(self) -> logging.Logger:
        return self.logger

    def config(self, dict_config=False, path='src/logging.yaml') -> logging.Logger:
        logger = logging.getLogger(self.tag)
        logger.setLevel(self.base_level)
        
        if dict_config:
            import yaml
            with open(path) as fd:
                conf = yaml.load(fd, Loader=yaml.SafeLoader)
            logging.config.dictConfig(conf['logging'])
        else:
            h = handler.FluentHandler(
                self.tag, 
                host=self.host, 
                port=self.port, 
                buffer_overflow_handler=self.overflow_handler
            )
            # Add a formatter to keep log metadata intact
            formatter = logging.Formatter('%(asctime)s - %(name)s - %(levelname)s - %(message)s')
            h.setFormatter(formatter)
            logger.addHandler(h)
        return logger

    def overflow_handler(self, pendings):
        unpacker = msgpack.Unpacker(BytesIO(pendings))
        for unpacked in unpacker:
            print("Buffer overflow data:", unpacked)

if __name__ == '__main__':
    local_host = 'localhost' # Use your Fluentd container IP or service name if needed
    el = ElasticLogger('test-logger', local_host, 24224, logging.NOTSET, False).build()
    # Send structured data as a dictionary (no stringified JSON!)
    el.error({"json": "message", "log": "work dammit"})
    # You can also send plain text logs directly
    el.error("This is a plain text error message")

3. Verification Steps

After making these changes:

  1. Check Docker Network Connectivity: Run docker exec <your-fluentd-container-name> ping elasticsearch to confirm Fluentd can reach Elasticsearch.
  2. Check Elasticsearch Indices: Run curl http://localhost:9200/_cat/indices?v—you should see an index named chris.risley-YYYYMMDD.
  3. Configure Kibana: Create an index pattern matching chris.risley-* in Kibana's Management tab to view your logs.
  4. Inspect Fluentd Logs: Run docker logs <your-fluentd-container-name> to confirm no new connection or parsing errors.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 08:32:33