Flink 正在从结构化流处理走向多模态数据底座

2026-07-08 29 预计阅读时间: 1 分钟
来源: my.oschina.net AI 摘要 Original link

Disclaimer: This article is an AI-assisted summary. Read it together with the original source when precision matters. The summary may omit context, version differences, or edge cases and is not official documentation.

预计阅读时间:12 分钟

过去的数据工程大多围绕结构化记录展开:订单、日志、点击、交易、指标、维表、事实表。数据被切成一行行 schema 清晰的记录,处理任务也相对明确:清洗、关联、聚合、统计。现在问题变了。业务链路里开始同时出现文本、图片、音频、视频、向量、模型推理结果和传统字段,Apache Flink 这样的流式计算系统,也需要从“处理表和字段”扩展到“承载持续产生的多模态信号”。

这不是说结构化处理过时了。恰恰相反,多模态数据要进入生产系统,仍然需要事件时间、状态、窗口、Exactly-once、连接器、反压控制这些老本领。变化在于:Flink 处理的不再只是 user_id + item_id + click_time,而可能是 user_id + image_url + embedding + model_score + event_time

结构化记录仍然是骨架

多模态听起来像是图片、音频、视频的天下,但在工程链路里,它通常不会脱离结构化元数据单独存在。

一张商品图需要商品 ID、类目、上传时间、审核状态;一段用户语音需要会话 ID、语言、设备、时间戳;一次模型推理需要输入版本、模型版本、特征版本和置信度。没有这些字段,多模态对象很难被检索、关联、治理和回放。

所以更现实的抽象不是“Flink 直接处理图片”,而是:

  • 用结构化事件描述多模态对象的位置、元数据和处理状态;
  • 用流处理编排抽取、推理、向量化、过滤、聚合等步骤;
  • 把大对象存放在对象存储、湖仓或专用服务里,Flink 在流上搬运引用、特征和结果;
  • 对低延迟场景,把模型推理或向量检索放进算子、异步 I/O 或外部服务调用中。

这让 Flink 的角色更像一条持续运行的数据生产线:它不一定亲自“吞下”所有图片和视频字节,但它负责让每个对象在正确的时间被处理、关联、更新和落库。

多模态流处理的关键变化

传统 Flink 作业常见输入是 Kafka 里的 JSON、Avro、Protobuf 或 CDC 变更日志。多模态链路会让事件变得更“胖”,也更不稳定。

典型变化包括:

  • 字段从标量扩展到引用和向量:事件中可能包含 image_uriaudio_uriembeddingcaptionocr_text 等字段。
  • 处理步骤从 SQL 聚合扩展到模型调用:清洗之后可能要调用 OCR、ASR、图像分类、文本 embedding 或 rerank 模型。
  • 延迟分布更宽:一次模型推理可能几十毫秒,也可能因为外部服务拥塞变成几秒。
  • 状态更复杂:同一个业务对象可能先到元数据,再到图片处理结果,再到人工审核结果,需要流式关联和补全。
  • 成本成为一等约束:模型推理、向量化和大对象读取都比普通字段计算贵,不能无脑全量重算。

这也是 Flink 仍然有价值的地方:它擅长处理无界数据流、维护状态、按事件时间计算,并且能把外部系统连接成一条可恢复的流水线。

下面是一个可改造的最小示例。假设图片本体存放在对象存储,Kafka 只传图片 URI 和业务元数据;另一个 Kafka topic 接收模型服务产出的标签和置信度。Flink SQL 负责把两条流按 image_id 关联,并写入下游结果 topic。

运行前需要替换:

  • Kafka 地址 localhost:9092
  • topic 名称;
  • Flink 发行版中对应的 Kafka SQL Connector jar。
-- image_events: 业务系统产生的图片事件
CREATE TABLE image_events (
  image_id STRING,
  user_id STRING,
  image_uri STRING,
  category STRING,
  event_time TIMESTAMP(3),
  WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND
) WITH (
  'connector' = 'kafka',
  'topic' = 'image-events',
  'properties.bootstrap.servers' = 'localhost:9092',
  'properties.group.id' = 'flink-image-events',
  'scan.startup.mode' = 'latest-offset',
  'format' = 'json'
);

-- image_model_results: 模型服务产生的识别结果
CREATE TABLE image_model_results (
  image_id STRING,
  model_name STRING,
  model_version STRING,
  label STRING,
  confidence DOUBLE,
  inference_time TIMESTAMP(3),
  WATERMARK FOR inference_time AS inference_time - INTERVAL '5' SECOND
) WITH (
  'connector' = 'kafka',
  'topic' = 'image-model-results',
  'properties.bootstrap.servers' = 'localhost:9092',
  'properties.group.id' = 'flink-image-model-results',
  'scan.startup.mode' = 'latest-offset',
  'format' = 'json'
);

-- 下游结果:结构化元数据 + 多模态模型结果
CREATE TABLE enriched_image_events (
  image_id STRING,
  user_id STRING,
  image_uri STRING,
  category STRING,
  label STRING,
  confidence DOUBLE,
  model_version STRING,
  event_time TIMESTAMP(3)
) WITH (
  'connector' = 'kafka',
  'topic' = 'enriched-image-events',
  'properties.bootstrap.servers' = 'localhost:9092',
  'format' = 'json'
);

INSERT INTO enriched_image_events
SELECT
  e.image_id,
  e.user_id,
  e.image_uri,
  e.category,
  r.label,
  r.confidence,
  r.model_version,
  e.event_time
FROM image_events e
JOIN image_model_results r
ON e.image_id = r.image_id
AND r.inference_time BETWEEN e.event_time - INTERVAL '1' MINUTE
                         AND e.event_time + INTERVAL '10' MINUTE
WHERE r.confidence >= 0.8;

这个例子没有把模型推理塞进 Flink SQL。它采用的是更稳妥的生产形态:Flink 负责编排和关联,模型服务独立扩缩容。对于 GPU 推理、批量 embedding、复杂重试,这种拆分通常更容易控制成本和故障边界。

如果你的场景需要在 Flink 作业里直接调用模型服务,可以这样实践一个异步 I/O 算子。下面是伪项目级示例,假设外部 HTTP 服务接收 image_uri 并返回标签。

# requirements: apache-flink, aiohttp
# 运行方式需按你的 Flink/PyFlink 环境调整:python async_inference_job.py

import json
import aiohttp
from pyflink.common import Types
from pyflink.datastream import StreamExecutionEnvironment
from pyflink.datastream.functions import MapFunction

INFERENCE_ENDPOINT = "http://localhost:8080/classify"

class ParseEvent(MapFunction):
    def map(self, value):
        event = json.loads(value)
        return json.dumps({
            "image_id": event["image_id"],
            "user_id": event["user_id"],
            "image_uri": event["image_uri"]
        })

# 说明:PyFlink 的异步 I/O API 会随版本变化。
# 这里给出可改造的作业骨架;生产环境建议使用 Flink 官方版本对应的 AsyncFunction 写法。
async def call_model(image_uri):
    async with aiohttp.ClientSession() as session:
        async with session.post(INFERENCE_ENDPOINT, json={"image_uri": image_uri}) as resp:
            resp.raise_for_status()
            return await resp.json()


def main():
    env = StreamExecutionEnvironment.get_execution_environment()
    env.set_parallelism(2)

    events = env.from_collection(
        collection=[
            '{"image_id":"img-1","user_id":"u-1","image_uri":"s3://bucket/a.jpg"}',
            '{"image_id":"img-2","user_id":"u-2","image_uri":"s3://bucket/b.jpg"}'
        ],
        type_info=Types.STRING()
    )

    parsed = events.map(ParseEvent(), output_type=Types.STRING())
    parsed.print()

    env.execute("multimodal-event-skeleton")

if __name__ == "__main__":
    main()

这段 Python 代码重点不是展示完整推理,而是给出作业边界:Flink 流里应该传事件、引用和结果;真正的模型调用需要结合你使用的 Flink 版本、异步 I/O API、超时、并发度和重试策略落地。

架构上要避开的坑

多模态数据流上生产,最容易踩的坑不是“不会写 SQL”,而是把数据、模型和计算边界揉成一团。

几个判断标准很实用:

  • 不要把大文件直接塞进 Kafka:传 URI、hash、大小、媒体类型和元数据,文件放对象存储或湖仓。
  • 模型调用要有超时和降级:外部推理服务会抖动,Flink 算子不能无限等待。
  • 向量字段要控制大小:embedding 可以进入流,但高维向量会放大网络、状态和 checkpoint 成本。
  • 状态 TTL 必须明确:多模态结果可能晚到,窗口要给足容忍度,但不能让状态无限增长。
  • 重算路径要提前设计:模型升级后是否全量重跑?只重跑低置信度样本?还是按版本并存?这些会影响 topic、表结构和存储布局。

Flink 的优势是持续处理和状态一致性,不是替代所有媒体处理系统。图片解码、视频抽帧、GPU 推理、向量索引,仍然可以交给更专门的服务;Flink 负责把它们接入一条可观测、可恢复、可演进的数据流。

落地检查清单

准备把 Flink 用作多模态流式底座时,可以先问五个问题:

  1. 事件 schema 是否同时描述了业务对象、媒体引用、模型版本和处理状态?
  2. 大对象是否从流里剥离,只在事件中传可追踪 URI 和校验信息?
  3. 模型推理是放在 Flink 内部异步调用,还是拆成独立服务和结果流?
  4. checkpoint、状态大小、向量字段和窗口等待时间是否经过压测?
  5. 模型升级、数据回放、低置信度样本重处理有没有明确路径?

从结构化到多模态,Flink 的核心价值没有变:它仍然是在无界事件上做可靠计算。真正变的是事件内容更丰富、处理链路更长、外部系统更多。把结构化元数据当骨架,把多模态对象当被编排的生产资料,Flink 才能稳稳地站在新一代数据处理链路的中间。


相关推荐