背景游戏集群通常由大量分散的游戏节点构成多台服务器各自独立输出玩家行为日志。如果登录每台机器翻日志排查问题效率极低。我们需要一套统一的日志采集平台将所有游戏节点日志集中收集、实时入库、支持检索并自动清理过期数据。传统方案会选用 Hadoop 大数据栈但对于较小规模的游戏运营日志场景过于笨重。本文落地一套游戏运营日志统一平台采用 Filebeat Kafka ClickHouse对比传统 Hadoop 架构拆解完整数据流原理重点讲解 Filebeat 日志采集方案选型。一、架构选型Hadoop vs KFC传统 Hadoop 生态HDFSSparkHive适合 PB 级、长期归档、大规模离线分析场景游戏运营日志场景明显过重组件繁多Zookeeper、HDFS、YARN、Hive、Spark集群部署调试耗时任意组件故障都会导致整个链路瘫痪资源开销巨大至少需要多台机器组成集群常驻占用大量 CPU / 内存游戏中小规模日志场景资源严重浪费实时性差以离线批处理为主日志入库延迟分钟小时运营查玩家实时行为日志无法满足运维门槛高需要专职大数据运维版本升级、集群容错、数据修复成本高。FilebeatKafkaClickHouse KFC轻量架构游戏运营日志场景更合适组件极简全部容器化部署单台云服务器即可承载后端所有节点上报的日志分布式采集游戏服务器分散在多台机器每台节点部署 Filebeat 作为采集 Agent自动监控本地日志文件把日志统一上报到 Kafka实现多节点日志集中管理流式实时链路日志产生后秒级入库运营后台可以立刻检索玩家行为列式存储查询快千万级日志下按玩家 ID、模块、时间筛选毫秒返回自带 TTL 自动清理按月自动删除过期游戏日志不用写定时清理脚本削峰容错Kafka 做缓冲游戏服务器瞬时日志峰值不会打崩存储层。选型结论 海量离线归档、深度数据挖掘 → Hadoop/Spark 多节点游戏集群实时运营日志检索、短期存储 → FilebeatKafkaClickHouse二、全链路原理数据流链路多台游戏节点输出单行 JSON 日志 → 每台节点部署 Filebeat 采集日志文件直接解析 JSON 字段→ 统一推送到 Kafka Topic → ClickHouse Kafka 引擎外表消费消息 → 物化视图自动写入 MergeTree 主表 → Java 后台查询 ClickHouse 做运营检索。各组件职责Filebeat分布式日志采集 JSON 预处理多台游戏节点各自部署 Filebeat Agent监控本机日志文件增量。核心能力直接解析单行 JSON把 JSON 内的 key-value 提取出来不再把整条 JSON 塞到一个 message 字段。重点旧版本 Filebeat 使用type:log采集存在不少硬伤下面单独对比原理。Kafka消息缓冲层解耦采集和入库。游戏高峰期日志突增时Kafka 缓存消息削峰ClickHouse 消费速度慢也不会丢失日志支持多消费组后续可以新增其他消费端如日志告警而不影响现有 ClickHouse 链路。ClickHouse存储与查询层Kafka 引擎外表只负责订阅 Kafka 主题不持久化数据物化视图监听外表流入数据自动把数据写入主表MergeTree 主表列式持久化存储支持分区、TTL、高性能查询。三、Filebeat 采集方案选型type:log 旧方案问题 filestreamndjson 原理旧方案type: log采集模式不推荐早期 Filebeat 采集配置使用type: log读取日志文件该模式存在几个在游戏多节点场景下非常突出的问题文件状态管理缺陷type:log基于原生 input 实现依赖文件 inode 路径记录读取位点。游戏服务日志滚动logrotate时很容易出现位点错乱重复消费日志或者漏采日志。游戏集群节点多一旦发生日志丢失 / 重复排查成本极高。多行合并逻辑容易干扰 JSON 单行日志type:log内置多行合并处理器默认行为容易误将多条单行 JSON 合并成一条非法 JSON导致下游 Kafka、ClickHouse 解析失败。虽然可以关闭多行但配置繁琐。性能与资源稳定性差type:log是老实现在高并发日志写入场景下文件监听 CPU 占用更高。多游戏节点同时采集大量 Agent 并发时更容易出现卡顿。JSON 解析嵌套能力弱即便搭配decode_json处理器也是读取完整一行放入message字段之后再二次解析 JSON。解析失败时错误信息简陋游戏日志偶尔出现脏 JSON很难定位脏数据来源。旧方案完整链路读取日志行 → 存入message字段 → decode_json 处理器二次解析 → 输出到 Kafka。 缺点二次解析多一层开销解析失败直接整条消息变成 message 原始字符串污染下游表结构。新方案type: filestream ndjsonfilestream 是 Filebeat 新一代文件采集输入专门替换老旧loginput针对日志滚动、增量读取做了大量优化。 配合ndjson解析器直接在采集阶段解析单行 JSON。 示例原始日志行{time:1789714834611,date:2026-09-18 15:00:34,sid:16,playerId:yyy1013,module:net,ext:{model:logout,type:-1}}解析之后推送到 Kafka 的消息结构{ time:1789714834611, date:2026-09-18 15:00:34, sid:16, playerId:yyy1013, module:net, ext:{model:logout,type:-1} }链路读取日志行 → ndjson 直接解析 JSON 键值 → 输出顶层字段。 优点filestream 使用文件偏移量记录位点日志滚动稳定性强适合多游戏节点长期运行ndjson 原生解析解析失败单独标记 error_key不会污染整条数据没有 message 包裹层减少下游 ClickHouse JSON 解析压力。踩坑Filebeat 默认会追加 beat 元数据最开始直接使用 Filebeat 采集输出到 Kafka 的数据会自带timestamp、metadata、beat这些元字段{timestamp:2026-09-18T10:21:02.754Z,metadata:{beat:filebeat},time:1789714834611,ext:{model:logout}}这些字段不是业务日志会污染 ClickHouse 表。✅ 解决在 filebeat 配置增加processors删除不需要的字段processors: - drop_fields: fields: [timestamp, metadata, beat, host]Filebeat 最终可用配置 filebeat.yml部署在每一台游戏节点filebeat.inputs: - type: filestream paths: - /game/logs/*.log parsers: - ndjson: overwrite_keys: true add_error_key: true processors: - drop_fields: fields: [timestamp, metadata, beat, host] output.kafka: hosts: [kafka:9092] topic: item_json_topic partition.round_robin: reachable_only: false required_acks: 1 compression: gzip max_message_bytes: 1000000四、ClickHouse 全套建表语句游戏运营日志三张表主表 (存储)、Kafka 外表 (消费)、物化视图 (自动同步)1. 业务主表 item_logMergeTree带 TTL 1 个月自动删除CREATE TABLE IF NOT EXISTS item_log ( time UInt64 COMMENT 13位毫秒时间戳, date String COMMENT 日志时间字符串, playerId String COMMENT 玩家ID, sid String COMMENT 会话ID, reason String COMMENT 操作原因, ext String COMMENT 扩展嵌套JSON字符串, module String COMMENT 游戏模块 ) ENGINE MergeTree() ORDER BY (playerId, time) PARTITION BY toYYYYMM(toDateTime(time / 1000)) TTL toDateTime(time / 1000) INTERVAL 1 MONTH DELETE COMMENT 游戏运营日志主表数据保存1个月自动清理;说明time/1000毫秒转秒用于 DateTime TTL 计算ORDER BY (playerId, time)游戏业务最常用查询维度按玩家 ID 时间检索ext 字段为 String保存原始嵌套 JSON查询嵌套字段用 JSONExtract 函数。2. Kafka 引擎外表 item_log_kafkaCREATE TABLE IF NOT EXISTS item_log_kafka ( time UInt64, date String, playerId String, sid String, reason String, ext String, module String ) ENGINE Kafka SETTINGS kafka_broker_list kafka:9092, kafka_topic_list item_json_topic, kafka_group_name clickhouse-game-log-group, kafka_format JSONEachRow, kafka_skip_broken_messages 1;kafka_skip_broken_messages 1遇到格式错误日志跳过防止消费链路中断游戏日志偶尔会有脏数据。3. 物化视图 item_log_mv自动把 Kafka 外表数据写入主表CREATE MATERIALIZED VIEW IF NOT EXISTS item_log_mv TO item_log AS SELECT time, date, playerId, sid, reason, ext, module FROM item_log_kafka;创建物化视图后自动监听 Kafka 流入消息无需任何定时任务实时写入主表。五、docker部署ClickHouse docker-compose.ymlservices: clickhouse: container_name: clickhouse image: clickhouse/clickhouse-server:26.7.3 ports: - 8123:8123 - 9000:9000 environment: - CLICKHOUSE_USERadmin - CLICKHOUSE_PASSWORD123456 - CLICKHOUSE_DEFAULT_ACCESS_MANAGEMENT1 volumes: - ./clickhouse_data:/var/lib/clickhouse ulimits: nofile: soft: 262144 hard: 262144启动命令docker-compose -f clickhouse-compose.yml up -d进入 clickhouse 客户端创建账号容器启动后执行一次docker exec -it clickhouse clickhouse-clientCREATE USER admin IDENTIFIED BY 123456; GRANT ALL ON *.* TO admin;kafka-compose.ymlversion: 3.8 services: zookeeper: image: confluentinc/cp-zookeeper:7.6.0 container_name: zookeeper environment: ZOOKEEPER_CLIENT_PORT: 2181 ZOOKEEPER_TICK_TIME: 2000 ports: - 2181:2181 kafka: image: confluentinc/cp-kafka:7.6.0 container_name: kafka depends_on: - zookeeper ports: - 9092:9092 - 9093:9093 environment: KAFKA_BROKER_ID: 1 KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181 KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1 # 两个监听器容器内不同端口解决端口冲突 KAFKA_LISTENERS: INTERNAL://0.0.0.0:9092,EXTERNAL://0.0.0.0:9093 KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: INTERNAL:PLAINTEXT,EXTERNAL:PLAINTEXT KAFKA_INTER_BROKER_LISTENER_NAME: INTERNAL # 广播地址 KAFKA_ADVERTISED_LISTENERS: INTERNAL://kafka:9092,EXTERNAL://[your ip]:9093 KAFKA_LOG4J_ROOT_LOGLEVEL: INFO restart: unless-stopped六、环境部署验证流程启动 Kafka 容器创建 topicitem_json_topic启动 ClickHouse 容器测试 JDBC 连接root/123456依次执行三张表建表 SQL在所有游戏节点部署 Filebeat配置上面 filebeat.yml 并启动往游戏日志文件写入测试 JSON 日志执行 count (*) 验证数据入库七、ext 动态扩展字段原理与 Java 业务层实现7.1 ext 字段设计原理游戏运营日志存在大量业务动态扩展字段不同业务模块、不同版本会新增自定义埋点字段如果全部映射为 ClickHouse 表物理列每次新增埋点就需要 DDL 改表运维成本极高。本方案采用Map(String,String)类型的ext字段存储所有动态业务扩展属性固定业务字段time、playerId、sid、module作为表真实物理列查询、排序、过滤性能高所有不确定、动态新增的埋点属性统一存入ext Map格式ext[key]读取对应 valueMap 字段支持按 key 做条件过滤、模糊匹配、范围查询数值类字段查询时使用toFloat64OrZero()做类型转换解决 Map 内部全部为字符串无法直接数值比较的问题优势业务新增埋点不需要修改 ClickHouse 表结构服务端直接写入即可对多游戏节点、频繁迭代的游戏业务非常友好注意Map 字段 key 查询性能弱于物理列高频过滤字段建议升级为物理列。参考日志如下{time:1789972731683,date:2026-09-21 14:38:51,sid:16,playerId:pd34,module:item,ext:{reason:10070,modelId:20000003,changeNum:500,fromNum:10200}} {time:1789972733560,date:2026-09-21 14:38:53,sid:16,playerId:pd34,module:item,ext:{reason:10070,modelId:20000003,changeNum:500,fromNum:9700}} {time:1789972735801,date:2026-09-21 14:38:55,sid:16,playerId:pd34,module:item,ext:{reason:10070,modelId:20000003,changeNum:500,fromNum:9200}} {time:1789972736636,date:2026-09-21 14:38:56,sid:16,playerId:pd34,module:quest,ext:{questId:10005}} {time:1789972736637,date:2026-09-21 14:38:56,sid:16,playerId:pd34,module:item,ext:{reason:10004,modelId:20000003,changeNum:3000,fromNum:8700}} {time:1789972936538,date:2026-09-21 14:42:16,sid:16,playerId:s12,module:item,ext:{reason:10039,modelId:10000002,changeNum:60,fromNum:0}} {time:1789972936538,date:2026-09-21 14:42:16,sid:16,playerId:s12,module:item,ext:{reason:10025,modelId:10000001,changeNum:4000,fromNum:5038}}ClickHouse 查询示例-- 查询ext里面level10的日志 select * from item_log where ext[level] 10; -- 对ext数值做范围查询必须做类型转换 select * from item_log where toFloat64OrZero(ext[cost]) between 100 and 1000; -- ext字段模糊匹配 select * from item_log where ext[itemName] like %钻石%;7.2 Java 业务层动态过滤实现业务接口需要同时支持固定字段过滤ext 动态扩展字段过滤。 约定入参规则前端传入字段名以ext.开头则代表这是 Map 内部动态 key解析后走ext[xxx]表达式其余字段直接映射数据库物理列。核心逻辑SQL 拼接区分普通物理字段与 ext 扩展字段数值条件强制增加类型转换同时使用 JDBC 预编译参数占位符?杜绝 SQL 注入风险。核心业务代码/** * 构建查询SQL支持固定字段 ext动态Map字段过滤 */ private QuerySpec buildQuerySpec(ReqQueryItemLog req) { StringBuilder whereSql new StringBuilder(BASE_TABLE_SQL); ListObject args new ArrayList(); // 基础固定条件时间、模块、玩家ID、服务器ID等物理列 if (req.getStartTime() 0) { whereSql.append( and time ?); args.add(req.getStartTime()); } if (req.getEndTime() 0) { whereSql.append( and time ?); args.add(req.getEndTime()); } if (StringUtils.isNotBlank(req.getModule())) { whereSql.append( and module ?); args.add(req.getModule().trim()); } if (StringUtils.isNotBlank(req.getPlayerId())) { whereSql.append( and playerId ?); args.add(req.getPlayerId().trim()); } // 解析动态过滤条件 appendDynamicFilters(whereSql, args, req); return new QuerySpec(whereSql.toString(), args); } /** * 分发过滤条件区分普通物理字段 和 ext扩展Map字段 * 约定field以 ext. 开头代表map内部key */ private void appendDynamicFilters(StringBuilder whereSql, ListObject args, ReqQueryItemLog req) { if (req.getFilters() null || req.getFilters().isEmpty()) { return; } for (ReqQueryItemLog.FilterItem item : req.getFilters()) { if (item null || StringUtils.isBlank(item.getField())) { continue; } String operator normalizeOperator(item.getOperator()); if (item.getField().startsWith(ext.)) { // 动态扩展字段走Map表达式 ext[key] appendExtCondition(whereSql, args, item, operator); } else { // 普通物理表字段 appendFixedCondition(whereSql, args, item, operator); } } } /** * 处理ext Map字段的条件、like、between、in * 数值范围查询使用 toFloat64OrZero() 做字符串转数值 */ private void appendExtCondition(StringBuilder whereSql, ListObject args, ReqQueryItemLog.FilterItem item, String operator) { String extKey sanitizeExtKey(item.getField().substring(4)); if (StringUtils.isBlank(extKey)) { return; } // ClickHouse Map取值表达式 ext[keyName] String extExpr ext[ extKey ]; switch (operator) { case between: // map内部存储都是字符串数值比较必须转换 String[] range parseRange(item.getRangeValue()); whereSql.append( and toFloat64OrZero().append(extExpr).append() between ? and ?); args.add(range[0]); args.add(range[1]); break; case in: appendInCondition(whereSql, args, extExpr, item.getMultiValue()); break; case like: whereSql.append( and ).append(extExpr).append( like ?); args.add(item.getValue().trim()); break; default: whereSql.append( and ).append(extExpr).append( ).append(operator).append( ?); args.add(item.getValue().trim()); break; } } /** * 普通物理字段条件拼接 */ private void appendFixedCondition(StringBuilder whereSql, ListObject args, ReqQueryItemLog.FilterItem item, String operator) { String field sanitizeFixedField(item.getField()); if (field null) return; switch (operator) { case between: String[] range parseRange(item.getRangeValue()); whereSql.append( and ).append(field).append( between ? and ?); args.add(range[0]); args.add(range[1]); break; case in: appendInCondition(whereSql, args, field, item.getMultiValue()); break; case like: whereSql.append( and ).append(field).append( like ?); args.add(item.getValue().trim()); break; default: whereSql.append( and ).append(field).append( ).append(operator).append( ?); args.add(item.getValue().trim()); break; } } /** * in 条件通用处理使用预编译占位符防止SQL注入 */ private void appendInCondition(StringBuilder whereSql, ListObject args, String fieldExpr, String multiValue) { String[] values parseMultiValue(multiValue); if (values.length 0) return; whereSql.append( and ).append(fieldExpr).append( in (); for (int i 0; i values.length; i) { if(i0) whereSql.append(, ); whereSql.append(?); args.add(values[i]); } whereSql.append()); }7.3 入参示例前端传动态过滤条件示例查询ext.level 80{ filters:[ { field:ext.level, operator:, value:80 } ] }实际运行界面如下八、总结这套游戏运营日志平台解决多游戏节点日志分散难以统一查看的痛点抛弃笨重 Hadoop 架构采用 Filebeat filestreamndjson 采集、Kafka 削峰、ClickHouse 列式存储。 相比老旧type:log方案filestream 采集稳定性更强更适合多台游戏节点长期持续采集日志ndjson 在采集层直接解析 JSON减少下游计算压力。ext 字段采用字符串存储方案平衡 ClickHouse 查询能力与 Java 业务代码兼容性。全套配置、SQL 可以保存环境故障或者服务器重装时一键复现整套日志平台。