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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 18:44:58