使用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
相关产品推荐
相关产品推荐

