Apache Beam为PCollection添加Schema时触发'int'无encode属性错误
解决Apache Beam PCollection Schema添加时的AttributeError问题
问题现象
为Apache Beam的PCollection添加Schema时,添加"schema writer2" >> beam.Map(lambda x:x).with_output_types(beamSQL_schema.properties)步骤后触发报错:
AttributeError: 'int' object has no attribute 'encode'
移除该步骤后Pipeline运行正常,能正确输出properties实例:
properties( uuid='bf1', read_timestamp='2024-03-31T16:23:07.938Z', source_timestamp='2024-03-31T16:23:07.938Z', object='public_dwh_entity_role_properties', read_method='postgresql-backfill', stream_name='projects/locations/streams', schema_key='ff24413219', sort_keys=[1138, ''], payload={'id': 35, 'public_identifier': '30ee', 'creation_timestamp': '2023-10-17T08:15:48.620Z', 'last_update_timestamp': '2023-10-17T08:15:48.620Z', 'entity_role_public_identifier': '3c53a', 'property_type': 'PERCENTAGE', 'property_value': '100.00' }, source_metadata={'schema': 'public', 'table': 'properties', 'is_deleted': False, 'change_type': 'INSERT', 'tx_id': None, 'lsn': '', 'primary_keys': [] } )
完整报错栈如下:
Traceback (most recent call last): File "apache_beam/runners/common.py", line 1418, in apache_beam.runners.common.DoFnRunner.process File "apache_beam/runners/common.py", line 624, in apache_beam.runners.common.SimpleInvoker.invoke_process File "apache_beam/runners/common.py", line 1582, in apache_beam.runners.common._OutputHandler.handle_process_outputs File "apache_beam/runners/common.py", line 1695, in apache_beam.runners.common._OutputHandler._write_value_to_tag File "apache_beam/runners/worker/operations.py", line 239, in apache_beam.runners.worker.operations.SingletonElementConsumerSet.receive File "apache_beam/runners/worker/operations.py", line 198, in apache_beam.runners.worker.operations.ConsumerSet.update_counters_start File "apache_beam/runners/worker/opcounters.py", line 213, in apache_beam.runners.worker.opcounters.OperationCounters.update_from File "apache_beam/runners/worker/opcounters.py", line 265, in apache_beam.runners.worker.opcounters.OperationCounters.do_sample File "apache_beam/coders/coder_impl.py", line 1495, in apache_beam.coders.coder_impl.WindowedValueCoderImpl.get_estimated_size_and_observables File "apache_beam/coders/coder_impl.py", line 1506, in apache_beam.coders.coder_impl.WindowedValueCoderImpl.get_estimated_size_and_observables File "apache_beam/coders/coder_impl.py", line 209, in apache_beam.coders.coder_impl.CoderImpl.get_estimated_size_and_observables File "apache_beam/coders/coder_impl.py", line 1584, in apache_beam.coders.coder_impl.LengthPrefixCoderImpl.estimate_size File "apache_beam/coders/coder_impl.py", line 248, in apache_beam.coders.coder_impl.StreamCoderImpl.estimate_size File "apache_beam/coders/coder_impl.py", line 1769, in apache_beam.coders.coder_impl.RowCoderImpl.encode_to_stream File "apache_beam/coders/coder_impl.py", line 1170, in apache_beam.coders.coder_impl.SequenceCoderImpl.encode_to_stream File "apache_beam/coders/coder_impl.py", line 270, in apache_beam.coders.coder_impl.CallbackCoderImpl.encode_to_stream File "/home/user/Code/dwh-datflow-pubsub-noti/.env/lib/python3.7/site-packages/apache_beam/coders/coders.py", line 429, in encode return value.encode('utf-8') AttributeError: 'int' object has no attribute 'encode'
错误原因
报错核心是数据类型与Schema定义不匹配:
- Schema中
sort_keys定义为Optional[Sequence[str]],但实际数据里是[1138, ''],其中1138是int类型; - Schema中
payload.id定义为Optional[str],但实际数据是35(int类型); - 当使用
with_output_types指定Schema后,Beam的RowCoder会严格按照Schema类型编码,尝试对int类型调用encode('utf-8')方法,触发AttributeError;不指定Schema时,Coder不会做严格类型校验,因此能正常运行。
解决方案
需让数据类型与Schema定义保持一致,以下两种方式二选一:
方式1:修正Schema定义,匹配实际数据类型
调整Schema中不匹配的字段类型:
from typing import Union, Optional, Sequence, NamedTuple class properties(NamedTuple): uuid: Optional[str] read_timestamp: Optional[str] source_timestamp: Optional[str] object: Optional[str] read_method: Optional[str] stream_name: Optional[str] schema_key: Optional[str] sort_keys: Optional[Sequence[Union[str, int]]] # 改为支持字符串和整数 payload: Optional[ NamedTuple( 'payload', id=Optional[int], # 改为整数类型 public_identifier=Optional[str], creation_timestamp=Optional[str], last_update_timestamp=Optional[str], entity_role_public_identifier=Optional[str], property_type=Optional[str], property_value=Optional[str], ) # type: ignore ] source_metadata: Optional[ NamedTuple( 'source_metadata', schema=Optional[str], table=Optional[str], is_deleted=Optional[bool], change_type=Optional[str], tx_id=Optional[str], lsn=Optional[str], primary_keys=Optional[Sequence[str]], ) # type: ignore ]
方式2:修正数据类型,匹配Schema定义
在转换为properties实例前,将不匹配的int类型转为字符串:
with beam.Pipeline(options=beam_options) as p: notifications = ( p | "read gcs noti" >> beam.io.ReadFromText('properties.jsonl') | "to json" >> beam.Map(json.loads) | "fix data types" >> beam.Map(lambda x: { **x, 'sort_keys': [str(item) for item in x['sort_keys']], 'payload': {**x['payload'], 'id': str(x['payload']['id'])} }) | "schema writer1" >> beam.Map(lambda x: beamSQL_schema.properties(**x)) | "schema writer2" >> beam.Map(lambda x:x).with_output_types(beamSQL_schema.properties) | beam.Map(print) )
验证
修改后重新运行Pipeline,with_output_types步骤可正常工作,不再触发AttributeError。
内容的提问来源于stack exchange,提问作者user23958383
相关产品推荐
相关产品推荐

