使用Composer 2.6.2和Airflow 2.5.3发布Pub/Sub消息时遇错误求助
Airflow发布Pub/Sub消息报错排查与解决
环境
- Composer 2.6.2
- Airflow 2.5.3
第一个错误:AttributeError: 'bytes' object has no attribute 'encode'
DAG运行失败,日志详情:
[2024-03-27, 06:25:37 UTC] {logging_mixin.py:137} INFO - data:b'DNB-MS-MY,2024-03-26,Location' [2024-03-27, 06:25:37 UTC] {taskinstance.py:1778} ERROR - Task failed with exception Traceback (most recent call last): File "/opt/python3.8/lib/python3.8/site-packages/airflow/operators/python.py", line 175, in execute return_value = self.execute_callable() File "/opt/python3.8/lib/python3.8/site-packages/airflow/operators/python.py", line 192, in execute_callable return self.python_callable(*self.op_args, **self.op_kwargs) File "/home/airflow/gcs/dags/pm_cm_dags/nim_batch_request_location_dag.py", line 148, in publish_message pubsub_message = {'data': base64.b64encode(data.encode()).decode()} AttributeError: 'bytes' object has no attribute 'encode'
报错代码片段:
def publish_message(environ, **kwargs): assert isinstance(data, bytes) print(f"data:{data}") pubsub_message = {'data': base64.b64encode(data.encode()).decode()} print(f"pubsub_message:{pubsub_message}") topic_name='' gcp_pubsub_hook.PubSubHook().publish(project_id=PROJECT_ID,topic=topic_name,messages=pubsub_message)
修改代码后出现的第二个错误:TypeError: string indices must be integers
将代码修改为pubsub_message = {'data': base64.b64encode(data).decode()}后,仍报错:
gcp_pubsub_hook.PubSubHook().publish(project_id=PROJECT_ID,topic=topic_name,messages=pubsub_message) File "/opt/python3.8/lib/python3.8/site-packages/airflow/providers/google/common/hooks/base_google.py", line 475, in inner_wrapper return func(self, *args, **kwargs) File "/opt/python3.8/lib/python3.8/site-packages/airflow/providers/google/cloud/hooks/pubsub.py", line 130, in publish self._validate_messages(messages) File "/opt/python3.8/lib/python3.8/site-packages/airflow/providers/google/cloud/hooks/pubsub.py", line 152, in _validate_messages if "data" in message and isinstance(message["data"], str): TypeError: string indices must be integers
解决方案
修正后的代码如下:
def publish_message(environ, **kwargs): assert isinstance(data, bytes) print(f"data:{data}") # 直接对bytes类型的data进行base64编码,再解码为字符串 pubsub_message = {'data': base64.b64encode(data).decode()} print(f"pubsub_message:{pubsub_message}") topic_name = '' # 替换为实际的Pub/Sub Topic名称 # messages参数传入包含消息字典的列表,而非单个字典 gcp_pubsub_hook.PubSubHook().publish( project_id=PROJECT_ID, topic=topic_name, messages=[pubsub_message] )
错误原因解析
- 第一个错误:
data已经是bytes对象,只有字符串(str)才有.encode()方法,bytes无需二次编码,直接传入base64.b64encode()即可完成编码。 - 第二个错误:Airflow的
PubSubHook.publish()方法的messages参数要求是消息字典的列表。当传入单个字典时,方法会将其视为可迭代对象,遍历字典的key(字符串类型),并尝试对字符串做索引操作,因此触发TypeError。
内容的提问来源于stack exchange,提问作者user23700007
相关产品推荐
相关产品推荐

