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

如何在Faust中不使用Record定义app.Table?

解决Faust动态Schema场景下Table.relative_to_field的耦合问题

问题描述

在使用Schema Registry配合Faust处理流数据时,由于数据Schema会动态变更,我选择避开faust.Record(避免每次Schema变更都修改实体类),直接处理字典格式的消息。但遇到限制:app.Table的relative_to_field方法要求传入FieldDescriptorT类型参数,该类型与faust.Record强绑定,无法像流的group_by那样用lambda或自定义函数提取字段。

解决方案

方案1:动态生成临时Record类获取FieldDescriptor

利用Faust的FieldDescriptor与Record字段的关联特性,从Schema Registry拉取Schema后动态生成临时Record类,再获取对应时间字段的描述符。这种方式符合Faust原生API设计,无需修改窗口逻辑。

import faust
from datetime import timedelta
from schema_registry.client import SchemaRegistryClient, schema
from schema_registry.serializers.faust import FaustSerializer

# 基础初始化配置
topic_name = "practice4"
subject_name = f"{topic_name}-value"
serializer_name = f"{topic_name}_serializer"
bootstrap_server = "192.168.59.100:30887"
sr_server = "http://localhost:8081"

client = SchemaRegistryClient({"url": sr_server})
topic_schema = client.get_schema(subject_name)
fp_avro_schema = schema.AvroSchema(topic_schema.schema.raw_schema)

avro_fp_serializer = FaustSerializer(client, serializer_name, fp_avro_schema)
faust.serializers.codecs.register(name=serializer_name, codec=avro_fp_serializer)

app = faust.App('sample_app', broker=bootstrap_server)
faust_topic = app.topic(topic_name, value_serializer=serializer_name)

# 根据Avro Schema动态生成临时Record类
def create_dynamic_record(avro_schema: schema.AvroSchema):
    field_defs = {}
    for field in avro_schema.fields:
        # 映射Avro类型到Python基础类型,复杂嵌套类型可按需扩展
        if field.type == "string":
            field_type = str
        elif field.type in ["int", "long"]:
            field_type = int
        elif field.type in ["float", "double"]:
            field_type = float
        elif field.type == "boolean":
            field_type = bool
        elif isinstance(field.type, schema.RecordSchema):
            field_type = create_dynamic_record(field.type)
        else:
            field_type = object
        field_defs[field.name] = faust.Field(type=field_type)
    return faust.Record(fields=field_defs, name="DynamicStreamRecord")

# 生成动态Record并获取时间字段的FieldDescriptor(假设时间字段名为event_time)
DynamicRecord = create_dynamic_record(fp_avro_schema)
time_field = DynamicRecord.__fields__["event_time"]

# 初始化Table并关联时间字段
count_table = app.Table(
    'count_table', default=int,
).hopping(
    timedelta(minutes=10),
    timedelta(minutes=5),
    expires=timedelta(minutes=10)
).relative_to_field(time_field)

@app.agent(faust_topic)
async def process_fp(fps):
    async for fp in fps.group_by(lambda fp: fp["job_id"], name=f"{subject_name}.job_id"):     
        count_table[fp["job_id"]] += 1
        print(fp)

方案2:绕过relative_to_field,手动指定事件时间更新窗口

不依赖relative_to_field,在处理每条消息时手动提取事件时间,调用Table更新方法时传入timestamp参数,让窗口基于事件时间而非处理时间计算。这种方式完全脱离faust.Record限制,灵活性更高。

import faust
from datetime import timedelta
from schema_registry.client import SchemaRegistryClient, schema
from schema_registry.serializers.faust import FaustSerializer

# 基础初始化配置(同方案1)
topic_name = "practice4"
subject_name = f"{topic_name}-value"
serializer_name = f"{topic_name}_serializer"
bootstrap_server = "192.168.59.100:30887"
sr_server = "http://localhost:8081"

client = SchemaRegistryClient({"url": sr_server})
topic_schema = client.get_schema(subject_name)
fp_avro_schema = schema.AvroSchema(topic_schema.schema.raw_schema)

avro_fp_serializer = FaustSerializer(client, serializer_name, fp_avro_schema)
faust.serializers.codecs.register(name=serializer_name, codec=avro_fp_serializer)

app = faust.App('sample_app', broker=bootstrap_server)
faust_topic = app.topic(topic_name, value_serializer=serializer_name)

# 初始化Table时不使用relative_to_field
count_table = app.Table(
    'count_table', default=int,
).hopping(
    timedelta(minutes=10),
    timedelta(minutes=5),
    expires=timedelta(minutes=10)
)

@app.agent(faust_topic)
async def process_fp(fps):
    async for fp in fps.group_by(lambda fp: fp["job_id"], name=f"{subject_name}.job_id"):
        # 提取事件时间并转换为秒级时间戳(需与Faust配置一致)
        event_time = fp["event_time"]
        event_ts = event_time.timestamp() if hasattr(event_time, "timestamp") else event_time / 1000
        
        # 基于事件时间更新Table窗口
        current_count = count_table.get(fp["job_id"], timestamp=event_ts)
        count_table.set_value(fp["job_id"], current_count + 1, timestamp=event_ts)
        print(fp)

方案对比

  • 动态Record类:贴合Faust原生设计,窗口逻辑无需调整,但需维护Avro类型到Python类型的映射,嵌套类型需额外处理。
  • 手动指定timestamp:完全脱离faust.Record耦合,适配任意动态Schema,但需自行处理时间字段的格式转换,确保与Faust的时间单位配置匹配。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 18:15:50