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

使用Go向Firehose写入JSON文件遇解析及格式异常问题咨询

错误原因

  • 首次报错:你通过os.ReadFile读取到的firehoseData已经是JSON文件的原始字节流,额外调用json.Marshal对字节数组做序列化操作时,Go会将[]byte类型序列化为base64编码的字符串,最终传入Firehose的是base64文本而非原始JSON,下游解析服务(通常是Elasticsearch)无法识别非JSON格式内容,因此抛出解析异常。
  • 二次报错:你的本地JSON文件是带缩进、换行的格式化多行JSON,而Firehose要求单条记录必须为无内部换行的单行有效JSON对象,不符合格式要求因此抛出malformed错误。

正确实现代码

import (
    "encoding/json"
    "os"
    "log"
    "fmt"
    "github.com/aws/aws-sdk-go/aws"
    "github.com/aws/aws-sdk-go/aws/session"
    "github.com/aws/aws-sdk-go/service/firehose"
)

func insertIntoFireHose(sess *session.Session, hoseName string) error {
    svc := firehose.New(sess, aws.NewConfig().WithRegion("us-east-1"))
    // 读取JSON文件
    rawBytes, err := os.ReadFile("/temp/test.json")
    if err != nil {
        return fmt.Errorf("read file failed: %w", err)
    }
    
    // 代码侧统一做序列化,确保输出单行无换行JSON,兼容任意格式的输入JSON
    var jsonObj map[string]interface{}
    if err = json.Unmarshal(rawBytes, &jsonObj); err != nil {
        return fmt.Errorf("parse raw json failed: %w", err)
    }
    recData, err := json.Marshal(jsonObj)
    if err != nil {
        return fmt.Errorf("marshal single line json failed: %w", err)
    }

    // 构造Firehose请求
    rec := &firehose.Record{
        Data: recData,
    }
    input := &firehose.PutRecordInput{
        DeliveryStreamName: aws.String(hoseName),
        Record: rec,
    }

    res, err := svc.PutRecord(input)
    if err != nil {
        return fmt.Errorf("put record to firehose failed: %w", err)
    }
    log.Println("put record success, response:", res)
    return nil
}

额外注意事项

  • 若下游服务要求每条JSON记录末尾带换行符做分隔,可在recData末尾追加[]byte("\n")后再传入Record。
  • 不要忽略任何错误处理,原代码中忽略os.ReadFile、json.Marshal的错误会导致问题排查难度大幅提升。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.25 05:15:04