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

使用Airflow的SqsPublishOperator触发Matillion ETL任务遇参数类型错误

解决Airflow中SqsPublishOperator发送消息时的“Invalid type for parameter”错误

问题根源

你遇到的错误是因为SqsPublishOperator的message_content参数要求传入字符串类型,但你直接传递了Python字典,AWS SQS API不接受非字符串格式的消息内容,因此触发类型校验错误。

修复方案

需要将字典格式的消息内容转换为JSON字符串,具体步骤:

  1. 导入json模块
  2. 使用json.dumps()将消息字典序列化为字符串

修改后的完整代码

import os
import json  # 新增导入
from airflow import DAG
from airflow.operators.dummy import DummyOperator
from airflow.providers.amazon.aws.operators.sqs import SqsPublishOperator
from common.util import configuration, alerting

dag_id =  os.path.basename(__file__).replace(".py", "")
config = configuration.get_config(f'{dag_id}.yml')

# SQS queue to publish to in AWS
sqs_queue_name = 'run-etl-jobs'

args = {
          "ENVIRONMENT":  "dev",
          "account_name": "mydevaccount",
          "group": "Matillion_ETL_DEV",
          "project": "General"
          }

with DAG(dag_id=dag_id, **config['airflow'],
        on_failure_callback = opsgenie_alert.create_alert) as dag:
        start_task = DummyOperator(task_id='start')
        end_task = DummyOperator(task_id='end')

    publish_to_queue = SqsPublishOperator(
    task_id="publish_to_queue",
    sqs_queue=sqs_queue_name,
    # 将字典转为JSON字符串
    message_content=json.dumps({
        "group":       args['group'],
        "project":     args['project'],
        "version":     "default",
        "environment": args['account_name'],
        "job":         "Matillion ETL Job",
        })
    )

    start_task >> publish_to_queue >> end_task

额外检查项

  • 确认sqs_queue参数的值:部分Airflow版本中,SqsPublishOperator可能需要传入SQS队列的URL或ARN,而不是队列名称。如果修复后仍有问题,尝试替换为队列的完整URL(例如https://sqs.<region>.amazonaws.com/<account-id>/run-etl-jobs)。
  • 验证Airflow的AWS连接配置是否正确,确保具备向目标SQS队列发送消息的权限。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 17:25:19