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

Logstash配置错误:Kafka转BigQuery同步失败求助

Logstash同步Kafka数据到BigQuery配置错误解决

核心配置错误修正

你的配置文件中google_bigquery输出插件的csv_schema字段语法错误,Logstash配置统一使用=>作为键值分隔符,而非>=,这是导致启动失败的直接原因。

修正后的完整配置

input {
  kafka {
        bootstrap_servers => ["addcn102:9092,addcn103:9092,addcn104:9092"]
        group_id => "logstash_kafka"
        auto_offset_reset => "earliest"
        consumer_threads => 3
        topics => ["houstonbq"]
        # 注意:你的Kafka数据带有_jsonparsefailure标签,说明原始数据不是标准JSON,若数据是Logstash Ruby格式,可改用rubydebug codec
        # codec => "rubydebug"
        codec => "json"
  }
}

output {
  google_bigquery {
        project_id => "newcar8891"
        dataset => "logstash"
        json_key_file => "/home/shurui/bin/newcar8891-013ef3777de7.json"
        error_directory => "/opt/module/bqerror"
        flush_interval_secs => 30  # 数值类型无需加引号
        csv_schema => "message:STRING"  # 修正为正确的=>分隔符
        # 可选:若需要自动创建表,添加以下配置
        # auto_create_table => true
  }
}

额外问题处理

从你提供的Kafka数据来看,存在_jsonparsefailure标签,说明当前codec => "json"无法正确解析数据:

  • 你的Kafka数据格式是Logstash的Ruby哈希格式(使用=>而非JSON的:),并非标准JSON。如果Kafka中存储的确实是这种格式,需要将input的codec改为rubydebug;
  • 若Kafka中应该是标准JSON,检查生产者是否正确输出JSON格式数据,避免格式不匹配导致解析失败。

验证步骤

  1. 替换修正后的配置文件;
  2. 重新执行启动命令:logstash-8.2.3/bin/logstash -f config/logstash-sample.conf;
  3. 观察日志是否再无配置错误,同时检查BigQuery中是否成功写入message字段数据。

内容的提问来源于stack exchange,提问作者Little white programmer

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 21:02:10