使用Airflow的SqsPublishOperator触发Matillion ETL任务遇参数类型错误
解决Airflow中SqsPublishOperator发送消息时的“Invalid type for parameter”错误
问题根源
你遇到的错误是因为SqsPublishOperator的message_content参数要求传入字符串类型,但你直接传递了Python字典,AWS SQS API不接受非字符串格式的消息内容,因此触发类型校验错误。
修复方案
需要将字典格式的消息内容转换为JSON字符串,具体步骤:
- 导入
json模块 - 使用
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
相关产品推荐
相关产品推荐

