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

使用Go语言Azure Functions存储JSON消息到Azure存储队列的重复数据问题

解决Azure Functions Custom Handler队列消息JSON数据重复问题

问题背景

使用Go结合Custom Handler开发Azure Functions时,往Storage Queue存入JSON格式消息后,消息的message字段与metadata字段出现重复数据;存入非JSON格式消息(如纯文本)时无此问题。

原因分析

Azure Functions的Custom Handler队列触发器会自动识别JSON格式的消息,将消息中的顶层键值对提取并注入到metadata中,同时在message字段保留完整的原始JSON字符串,最终导致数据重复。

解决方案

方案一:包装JSON为单一顶层字段结构

将原始JSON消息放入一个单一字段的JSON对象中,避免触发器解析原始JSON的顶层字段到metadata。

修改入队函数(function1):

import "encoding/json"

func function1(c *gin.Context) {
    log.Printf("Start enqueueMessage")

    // 原始JSON数据
    rawJSON := "{\"aaa\" : \"aaa\"}"
    // 定义包装结构体
    type WrappedPayload struct {
        Payload string `json:"payload"`
    }
    wrapped := WrappedPayload{Payload: rawJSON}
    // 序列化为JSON字节
    wrappedBytes, err := json.Marshal(wrapped)
    if err != nil {
        log.Fatal("Error marshaling wrapped payload: ", err)
    }
    // base64编码
    base64BodyString := base64.StdEncoding.EncodeToString(wrappedBytes)

    _ulr, err := url.Parse(fmt.Sprintf("https://%s.queue.core.windows.net/%s", accountName, queueName))
    if err != nil {
        log.Fatal("Error parsing url: ", err)
    }

    credential, err := azqueue.NewSharedKeyCredential(accountName, accountKey)
    if err != nil {
        log.Fatal("Error creating shared key credential: ", err)
    }

    queueUrl := azqueue.NewQueueURL(*_ulr, azqueue.NewPipeline(credential, azqueue.PipelineOptions{}))
    ctx := context.TODO()

    messageUrl := queueUrl.NewMessagesURL()
    _, err = messageUrl.Enqueue(ctx, base64BodyString, 0, 0)
    if err != nil {
        log.Fatal("Error enqueueing message: ", err)
    }
    log.Printf("Message enqueued successfully")
}

修改消费函数(function2)以解析包装后的消息:

import "encoding/json"

func function2(c *gin.Context) {
    log.Printf("Start function")

    type QueueMessage struct {
        Data struct {
            MyQueueItem string `json:"myQueueItem"`
        } `json:"data"`
        Metadata struct {
            DequeueCount    string `json:"DequeueCount"`
            ExpirationTime  string `json:"ExpirationTime"`
            Id              string `json:"Id"`
            InsertionTime   string `json:"InsertionTime"`
            NextVisibleTime string `json:"NextVisibleTime"`
            PopReceipt      string `json:"PopReceipt"`
            Sys             struct {
                MethodName string `json:"MethodName"`
                RandGuid   string `json:"RandGuid"`
                UtcNow     string `json:"UtcNow"`
            } `json:"sys"`
            Payload string `json:"payload"` // 对应包装后的字段
        }
    }

    var queueMessage QueueMessage
    err := c.BindJSON(&queueMessage)
    if err != nil {
        log.Printf("Failed to bind request body: %v", err)
        return
    }

    // 解析原始JSON数据
    var rawData map[string]string
    err = json.Unmarshal([]byte(queueMessage.Metadata.Payload), &rawData)
    if err == nil {
        log.Printf("Raw message data: %v", rawData)
    }

    log.Printf("Queue message : %v", queueMessage)
}

方案二:将JSON转义为纯文本字符串

将原始JSON字符串转义为纯文本字面量,让触发器识别为普通文本而非JSON对象,避免字段解析。

修改入队函数(function1):

import "strconv"

func function1(c *gin.Context) {
    log.Printf("Start enqueueMessage")

    // 原始JSON数据
    rawJSON := "{\"aaa\" : \"aaa\"}"
    // 转义为字符串字面量(添加外层引号)
    escapedJSON := strconv.Quote(rawJSON)
    // base64编码
    base64BodyString := base64.StdEncoding.EncodeToString([]byte(escapedJSON))

    // 后续URL解析、凭证创建、入队逻辑不变
    _ulr, err := url.Parse(fmt.Sprintf("https://%s.queue.core.windows.net/%s", accountName, queueName))
    if err != nil {
        log.Fatal("Error parsing url: ", err)
    }

    credential, err := azqueue.NewSharedKeyCredential(accountName, accountKey)
    if err != nil {
        log.Fatal("Error creating shared key credential: ", err)
    }

    queueUrl := azqueue.NewQueueURL(*_ulr, azqueue.NewPipeline(credential, azqueue.PipelineOptions{}))
    ctx := context.TODO()

    messageUrl := queueUrl.NewMessagesURL()
    _, err = messageUrl.Enqueue(ctx, base64BodyString, 0, 0)
    if err != nil {
        log.Fatal("Error enqueueing message: ", err)
    }
    log.Printf("Message enqueued successfully")
}

修改消费函数(function2)以还原原始JSON:

import (
    "strconv"
    "encoding/json"
)

func function2(c *gin.Context) {
    log.Printf("Start function")

    type QueueMessage struct {
        Data struct {
            MyQueueItem string `json:"myQueueItem"`
        } `json:"data"`
        Metadata map[string]interface{}
    }

    var queueMessage QueueMessage
    err := c.BindJSON(&queueMessage)
    if err != nil {
        log.Printf("Failed to bind request body: %v", err)
        return
    }

    // 还原转义的JSON字符串
    rawJSON, err := strconv.Unquote(queueMessage.Data.MyQueueItem)
    if err == nil {
        // 解析原始JSON
        var rawData map[string]string
        json.Unmarshal([]byte(rawJSON), &rawData)
        log.Printf("Raw message data: %v", rawData)
    }

    log.Printf("Queue message : %v", queueMessage)
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 09:44:57