ByteDance、ファイル即時アップロードによる多モーダルベクトル検索の実践を共有
本文の状態
日本語全文を表示中
詳細モードで約20分の本文を読めます。
同じ出来事の情報源
この情報源を基点に整理
ByteDance Engineering
ByteDance は Flink と VikingDB を活用し、ファイルアップロードから即座に検索可能なリアルタイム多モーダルベクトルリンクを構築する実装手法と SQL 参照コードを発表した。
AI深層分析を開く2026年8月4日 20:58
AI深層分析
キーポイント
リアルタイムデータリンクの統合
TOS、Kafka、Flink、VikingDB を連携させることで、ファイルアップロードから数秒以内での検索可能化を実現し、従来のバッチ処理による遅延や整合性問題を解消する。
TOS-CDC による全量・増量のシームレス接続
Flink SQL Source Connector「TOS-CDC」が対象ストレージの全量スキャンと Kafka を介した增量イベント消費を自動で切り替え、データ漏れや重複を防ぐ。
VikingDB による自動ベクトル化とインデックス更新
標準的な画像・テキスト検索において VikingDB が内部モデルを用いて自動で読み込み、ベクトル生成、および索引更新を行うため、外部の GPU リソース管理が不要になる。
Flink 2.2 AI SQL を用いたカスタムモデル連携
Flink 2.2 の「CREATE MODEL」と「ML_PREDICT」機能により、SQL 内で直接方舟(Ark)のモデルを呼び出してベクトル化を行い、柔軟な多モーダル処理が可能となる。
TOS-CDC と Kafka の連携設定
イベント通知で Bucket 内のオブジェクト変更を検知し、Kafka にプッシュして TOS-CDC が継続的に感知する。
重要な引用
「ファイルがすでにアップロードされたこと」は「コンテンツがすでに AI で利用可能であること」とは等しくない
TOS-CDC は全量スキャン開始時間を記録し、その時点の Kafka 位置から增量イベントを消費することで、スキャン期間中のデータ変化もカバーする
Flink 2.2 AI SQL は現在邀測中であり、SQL 内でモデル呼び出しとベクトル書き込みを統一して管理できる
Topic Retention 大于“全量扫描最大耗时 + rewind.offset”,避免全量扫描结束时需要回拨的事件已经过期。
編集コメントを表示
編集コメント
本記事は、AI アプリケーションの生産環境におけるデータ遅延問題を解決するための具体的な技術実装例を提供しており、特にベクトルデータベースとストリーム処理の連携において参考となる。ただし、Flink 2.2 の AI SQL 機能は現在邀測中である点に留意し、導入を検討する際は公式情報の確認が必要である。
Source Article
元記事を日本語で読む
本文に関係しない購読案内、埋め込み通知、サイト内プロモーションは除いています。
オリジナル:Viking
2026年8月4日 19時12分 北京
企業が AI アプリケーションを概念実証(PoC)から本番環境へ移行する際、必ず遭遇するのが以下のようなシナリオです。商品画像のデータベースには毎日数千枚の新規写真が追加され、企業向け知識ベースには新たなドキュメントが絶えず流入し、訓練データプラットフォームでは新しいサンプルを即座に検知する必要があります。こうしたデータの最初の保管先は、通常オブジェクトストレージとなります。
しかし、「ファイルをアップロードした」ことと「AI がその内容を活用できる」ことは別問題です。「保存する」状態から「利用する」状態へ移行するには、一般的に以下のフローを経由します。
新規または変更されたオブジェクトの検出 → オブジェクトまたはメタデータの読み取り → モデル入力のクリーニングと組み立て → Embedding(ベクトル化)の生成 → ベクトルデータベースへの書き込み → 検索インデックスの更新
もし上記のプロセスが複数の定期タスクやスクリプトを繋ぎ合わせたものであれば、本番環境では以下の3 つのボトルネックに直面することが多くなります。
データ更新の遅延:新ファイルは次のスキャンサイクルまで待たされるため、検索対象となる情報が数時間、あるいはそれ以上遅れる可能性があります。
既存データと新規データの連携困難:全量スキャン中に新たなファイルがアップロードされると、インクリメンタル(増分)処理への切り替え時に、データの抜けや重複が発生しやすくなります。
モデルとベクトル書き込みの複雑なリンク構築:画像の取得、モデルサービスの提供、GPU リソースの確保、失敗時の再試行、ベクトルの書き込み、インデックスの更新など、それぞれを個別に実装する必要があります。
これらの課題を解決するため、本稿では火山引擎(Volcengine)の Flink と VikingDB を組み合わせて、「ファイルアップロード後、数秒で検索可能になる」リアルタイムなマルチモーダルデータリンクを構築する方法を紹介し、2 つの完全な Flink SQL 参考実装案を提示します。
一、一つのリンクですべての工程を統合
Flink と VikingDB の連携ソリューションは、上記のように分散していた工程を、継続して実行されるリアルタイムデータリンクへと集約します。
このデータ処理パイプラインを構成する主要なコンポーネントは以下の通りです。
TOS(オブジェクトストレージ): 画像、動画、テキスト、ドキュメント、および業務メタデータを格納します。多模态データの保管場所となります。
Kafka(メッセージキュー): TOS での PUT や DELETE などのオブジェクトイベントを処理します。ファイルに変更が生じると、自動的に Kafka にメッセージが送信されます。
Flink(ストリーム計算エンジン): 全量および増分データの取り込み、データクリーニング、ルーティング、モデル呼び出し、そして障害復旧を担当します。このパイプライン全体を制御する「オーケストレーション・エンジン」です。
VikingDB(ベクトルデータベース): データのベクトル化、保存、インデックス作成、およびオンライン検索を担います。ベクトルデータの最終的な保管先であり、検索サービスのバックエンドとなります。
方舟(Ark): 標準モデルでは対応できないカスタム要件を持つシーン向けに、多模态な埋め込み(Embedding)機能を提供します。組み込みモデルが不十分な場合に、独自のモデルをここに接続して利用できます。
2. 3 つの主要機能でリアルタイム AI データパイプラインを実現
1. TOS-CDC:オブジェクトストレージを継続的に更新されるテーブルへ変換
TOS-CDC は、オブジェクトストレージ向けの Flink SQL Source Connector です。もともと散在していたオブジェクトストレージ内のファイル変更を、連続的に処理可能なデータストリームに変換します。
ジョブ起動後、以下の動作を行います:
- 全量スキャンの開始時間を記録する
- 指定された TOS バケットに対して全量スキャンを実行する
- 全量スキャン完了後は、デフォルトでその開始時間に対応する Kafka の位置からオブジェクトイベントを消費し始める
- チェックポイントにより、全量スキャンの進捗状況と Kafka の消費位点を保存する
全量スキャン時点から Kafka 上の増分データを継続して消費することで、スキャン期間中に発生したオブジェクトの変更もすべてカバーできます。下流システムで安定した主鍵を持つ Upsert(「存在すれば更新し、なければ挿入」)処理を行うことで、重複イベントが発生しても最終的には単一のレコードに収束します。
注:TOS-CDC が受け取るのはファイルの場所に関する情報であり、ファイルそのもののコンテンツは受け取りません。
- VikingDB:書き込み・ベクトル化・検索を一つのループで完結させる
標準的な画像やテキストの検索シナリオでは、VikingDB のテーブル上でフィールドの意味とベクトルモデルを宣言するだけで済みます。例えば、フィールドを image として宣言し、doubao-embedding-vision を設定すれば、VikingDB が自動的に画像の読み込み、ベクトル化、およびインデックス更新を担当します。
業務側で追加にモデルサービスの維持管理や GPU リソースプールの用意、ベクトル取り込みプログラムの構築は不要です。Flink が変化するデータを継続的に書き込み、VikingDB がそれを検索可能なベクトル資産として蓄積・定着させます。
VikingDB Connector は Flink の Changelog に基づいて Upsert と Delete を実行できます。デフォルトでは同期書き込みが設定されています:
'async' = 'false'
書き込み直後に高速な検索が必要な業務の場合、async=true を直接有効化することは推奨されません。非同期書き込みはスループット優先の設計であり、Collection(データ集合)と Index の可視性に遅延が生じる可能性があります。
- Flink 2.2 AI SQL:モデル呼び出しを SQL に組み込む
一部の業務では、Flink 側で独自に Embedding を処理する必要があります。典型的なシナリオは以下の通りです:
- 画像、タイトル、タグ、OCR テキストを組み合わせて多モーダル入力を作成する
- 指定された方舟モデルとそのバージョンを使用する
- モデル呼び出しの並列度、スループット、コストを一元的に制御する
- 一つの Embedding を VikingDB、Kafka、特徴量ストア、または訓練サンプルへ同時に書き込む
- Java や Python のジョブを再実装することなくモデルを切り替える
ストリーミング計算プラットフォーム Flink バージョン 2.2 では、CREATE MODEL で方舟モデルを宣言し、ML_PREDICT を SQL 内で使用することでリアルタイム推論を実行できます。モデルの接続、データ処理、ベクトル書き込みが一つの SQL ジョブで統括的にオーケストレーションされます。
注:現在の Flink バージョン 2.2 の AI SQL は招待制テスト中です。ご要望の場合は火山公式サイトまでお問い合わせください。
三.事前準備チェックリスト
リンク構築を開始する前に、以下のリソースを準備してください:
- 準備項目:役割 / 必須か
- TOS Bucket とイベント通知ルール:多モーダルオブジェクトの保存と、オブジェクト変更イベントの Kafka への転送 / 必須
- Kafka インスタンス、トピック、ユーザー:TOS の増分イベントを受け取り、TOS-CDC が継続的に消費する / 必須
- ストリーミング計算 Flink バージョン:TOS-CDC、データ処理、VikingDB Sink の実行環境 / 必須
- TOS-CDC 招待制テスト資格:全量スキャンと増分イベントを一体化した Source を利用可能にする / 必須
- Flink 2.2 招待制テスト資格:CREATE MODEL と ML_PREDICT を使用して方舟モデルを呼び出す / 方案二のみ必要
- VikingDB インスタンスと API Key:自動ベクトル化の実行、または Flink が生成したベクトルの保存 / 必須
- 方舟推論エンドポイントと API Key:カスタム多モーダル Embedding の実行 / 方案二のみ必要
注:両方のアプローチについては後ほど詳しく解説します。
ファイルアップロードで即座に検索可能|リアルタイム多モーダルベクトルリンクの実践
- TOS バケットの準備とイベント通知の有効化
TOS のイベント通知機能は、バケット内のオブジェクトに変更が発生した際に、その事象を Kafka へプッシュします。このイベントメッセージには、バケット名、オブジェクトキー、イベントの種類、発生時刻などの情報が含まれており、TOS-CDC はこれらを常時監視することで、新規作成、上書き、削除の各操作を検知します。
TOS コンソールで対象バケットにイベント通知ルールを作成する際は、以下の設定が必要です。
- 配信先を「メッセージキュー Kafka 版」に指定します。
- 少なくとも tos:ObjectCreated:* を購読してください。オブジェクトの削除処理も必要になった場合は、後から tos:ObjectRemoved:* も追加で購読します。
- バッファ範囲に合わせて Prefix と Suffix を設定し、イベントの対象範囲を TOS-CDC の bucket 設定と整合させます。
- 対象となる Kafka インスタンス、Topic、Kafka ユーザー、および権限付与ロールを選択します。
注:詳細な手順については、火山エンジンオブジェクトストレージのドキュメント「イベント通知を Kafka へプッシュする方法」をご参照ください。(リンク先:https://docs.volcengine.com/docs/6349/1817509?lang=zh)
- Kafka インスタンス、Topic、およびアクセス権限の準備
メッセージキュー Kafka 版において、インスタンス、Topic、およびアクセスユーザーを作成します。TOS のイベント通知ルールでは、Kafka インスタンス ID、Topic 名、ユーザー情報、IAM ロールを参照する必要があります。このロールには、TOS から Kafka へのイベント配信を許可するためのシステム事前定義ポリシー KafkaAccessForTOS をバインドしておく必要があります。
併せて、以下の点も確認してください。
- Flink リソースプールから Kafka Bootstrap Servers へアクセス可能であること。
- Kafka ユーザーと認証パラメータが Flink ジョブ内で使用可能であること。
- Topic の Retention(保持期間)を「全量スキャンの最大所要時間 + rewind.offset」よりも長く設定し、全量スキャン完了後に遡って処理が必要なイベントが期限切れになるのを防ぎます。
- 流式計算 Flink 版の有効化
火山エンジン流式計算 Flink 版を有効化し、プロジェクトとジョブ実行に必要なリソースプールを作成します。さらに、Kafka、VikingDB、および方舟サービスとのネットワーク接続を確立してください。
注:TOS-CDC のベータテスト参加資格や Flink 2.2 バージョンのベータテスト参加資格が必要な場合は、火山エンジンの公式サイトまでお問い合わせください。
- VikingDB と方舟リソースの準備
VikingDB を有効化し、データ面アドレスと API Key を用意します。方案一では、対象 Collection で使用する自動ベクトル化モデルの種類、バージョン、次元数を確認する必要があります。方案二では、方舟の推論アクセスポイントを作成し、モデル名、出力次元、API Key を準備します。
注:すべてのアクセス認証情報は、流式計算 Flink 版の暗号化変数または実行環境変数を通じて注入することを推奨します。SQL 内に直接記述するのは避けてください。
三、2 つの SQL 方案:用途に応じて選んでください
両方の方案は、同じ TOS-CDC データ取り込みリンクを共有しますが、「どちらがベクトル化処理を担当するか」で異なります。
- 方案一:VikingDB の自動ベクトル化。標準的な画像・テキスト検索向け。最もシンプルなアーキテクチャです。データを VikingDB に渡せば、同社が自動的に Embedding を生成します。
- 方案二:Flink AI SQL で方舟を呼び出し。画像とテキストの融合処理、モデルの自主制御、ベクトルの多様な下流利用向け。ユーザーがモデル、入力、出力を完全にコントロールできます。
- TOS-CDC ソーステーブルの作成
CREATE TABLE tos_object_events (
object_key STRING NOT NULL,
object_url STRING,
bucket_name STRING,
file_name STRING,
object_etag STRING,
object_size BIGINT,
mtime TIMESTAMP_LTZ(3),
event_time TIMESTAMP_LTZ(3),
record_origin STRING,
PRIMARY KEY (object_key) NOT ENFORCED
) WITH (
'connector' = 'tos-cdc',
'path' = 'tos://my-bucket/images01/, tos://my-bucket/images02/',
'properties.bootstrap.servers' = 'kafka.example:9092',
'properties.group.id' = 'ingest-cg',
'topic' = 'object-events',
'scan.startup.mode' = 'initial'
);
本番環境では Kafka インスタンスに応じて認証情報やネットワークパラメータを追加する必要があります。また、以下の点に注意してください。
Kafka Topic のメッセージ保持期間は、全量スキャンから Kafka への切り替え完了まで、および設定された rewind.offset をカバーできる時間以上に設定してください。
増分コールバックのパラメータ rewind.offset はデフォルトで 0 ですが、必要に応じて変更可能です。また、位置が期限切れになった際にデータが静かに欠落するのを防ぐため、rewind.retention-miss-policy=fail を維持することを強く推奨します。
出力データの構築:
CREATE TEMPORARY VIEW image_put_events AS
SELECT
object_key AS id,
object_url AS image_uri,
object_etag,
COALESCE(event_time, mtime) AS update_time
FROM tos_object_events;
ここでは object_key を主鍵として使用しています。TOS-CDC と外部 Sink の双方が「少なくとも 1 回」のセマンティクス(At-Least-Once)を採用しているため、故障復旧時や全量・増分の接続時にレコードが重複して再プレイされる可能性があります。安定した主鍵を設定することで、VikingDB 内で重複する PUT 操作が Upsert として処理され、最終的に単一のデータに収束します。
- 案一:VikingDB の自動ベクトル化
まず VikingDB Catalog を作成します:
CREATE CATALOG viking WITH (
'type' = 'vikingdb',
'control-plane.host' = 'open.volcengineapi.com',
'region' = 'cn-beijing',
'project-name' = 'default',
'access-key' = '${secret_values.volc-ak}',
'secret-key' = '${secret_values.volc-sk}',
'data-plane.host' = '',
'api-key' = '${secret_values.vikingdb-api-key}'
);
次に、自動画像ベクトル化を有効にした Collection を作成します。
CREATE TABLE IF NOT EXISTS viking.default.realtime_image_assets (
id STRING,
image_uri STRING,
object_etag STRING,
update_time TIMESTAMP_LTZ(3),
PRIMARY KEY (id) NOT ENFORCED
) WITH (
'vikingdb.field.image_uri.type' = 'image',
'vikingdb.vectorize.dense.model-name' = 'doubao-embedding-vision',
'vikingdb.vectorize.dense.model-version' = '',
設定値として、vikingdb.vectorize.dense.dim に 2048 を、vikingdb.vectorize.dense.image-field に image_uri を指定します。
次に VikingDB への書き込みを行います。
INSERT INTO `viking`.`default`.`realtime_image_assets`
SELECT id, image_uri, object_etag, update_time
FROM image_put_events;この処理フローで実現できるのは、ジョブ初回起動時に指定したストレージバケット内の既存オブジェクトを一括インポートすることです。 (原文の技術表記: viking、default、realtime_image_assets)
新对象上传後はリアルタイムで書き込みを行います。
同じパスのオブジェクトが上書きされた場合は、安定した主キーに基づいて更新します。
ジョブ失敗時にはチェックポイントから復元し、Upsert によってイベントの再実行を防ぎます。
- 案二:Flink AI SQL を用いた方舟 Embedding の呼び出し
業務でカスタム多模態入力が必要、または特定のモデルを指定する必要がある場合、同一の TOS-CDC ソーステーブルを再利用できます。
まず SQL で方舟モデルを宣言します:
CREATE MODEL ark_multimodal_embedding
INPUT (payload STRING)
OUTPUT (embedding ARRAY)
WITH (
'provider' = 'ark',
'endpoint' = 'https://ark.cn-beijing.volces.com/api/v3/embeddings/multimodal',
'api-key' = '${secret_values.ark-api-key}',
'model' = 'doubao-embedding-vision-251215',
'model.dimensions'= '2048'
);
オブジェクト情報をアーク(Ark)の多モーダル入力形式に組み立てます:
CREATE TEMPORARY VIEW multimodal_payload AS
SELECT
id,
image_uri,
update_time,
CAST(
JSON_ARRAY(
JSON_OBJECT(
'type' VALUE 'image_url',
'image_url' VALUE JSON_OBJECT(
'url' VALUE CONCAT(
'https:///',
object_key
)
)
)
) AS STRING
) AS payload
FROM (
SELECT
object_key AS id,
object_url AS image_uri,
COALESCE(event_time, mtime) AS update_time,
object_key
FROM tos_object_events
) AS source_events;
例として、HTTPS でアクセス可能な画像 URL をモデルの入力に使用しています。本番環境では、アークサービスがそのアドレスを安全にアクセスできることを保証する必要があります。プライベートなバケットを使用する場合は、制御された一時署名付き URL や社内認証リンクを利用し、モデル呼び出しのためにバケット全体を公開読み取り用に設定することは推奨されません。
明示的なベクトルを保存するための VikingDB Collection を作成します:
CREATE TABLE IF NOT EXISTS viking.default.realtime_multimodal_assets (
id STRING,
image_uri STRING,
update_time TIMESTAMP_LTZ(3),
mixed_embedding ARRAY,
PRIMARY KEY (id) NOT ENFORCED
) WITH (
'vikingdb.field.mixed_embedding.type'= 'vector',
'vikingdb.field.mixed_embedding.dim' = '2048'
);
モデルを呼び出して VikingDB に書き込みます:
INSERT INTO viking.default.realtime_multimodal_assets
SELECT
id,
image_uri,
update_time,
embedding AS mixed_embedding
FROM ML_PREDICT(
TABLE multimodal_payload,
MODEL ark_multimodal_embedding,
DESCRIPTOR(payload)
);
もし Embedding をリアルタイム特徴量や訓練サンプル、メッセージ購読に利用する必要がある場合、EXECUTE STATEMENT SET を使用して複数の Sink を追加し、下流システムで同じモデル計算結果を共有できるようにします。
五、このルートの検証方法
1. Flink タスクが実行状態にあるか確認
Flink UI を通じて、既存の画像や動画ファイルが VikingDB に正常にインポートされていることを確認してください。データ量と TOS の整合性が取れているかも併せてチェックします。
2. ベクトルと検索結果の検証
VikingDB コンソールの「データセット → データプレビュー」で、TOS のパスに沿ってクエリを実行し、データが正しくデータセットに書き込まれているか確認します。
VikingDB の Collection 内のフィールドタイプやベクトル次元数を確認した上で、類似画像や関連テキストを用いて検索を実行し、直近でアップロードされたオブジェクトが正しく召回(検索結果にヒット)されるか検証します。以下の図のように、「小松鼠」と入力すると、関連する類似写真が検索結果として表示されます。
3. エンドツーエンドの遅延測定
以下の項目をそれぞれ記録します。
- TOS オブジェクトのアップロード時間
- Kafka イベント発生時刻
- Flink の処理時間
- VikingDB への書き込みが可視化されるまでの時間
- そのオブジェクトを初めて検索できるまでの時間
P50 および P95 の遅延を、実際のデータ規模と並行処理条件に基づいて評価した上で、Flink の並行数、モデルのスループット、Sink Flush Interval、そして VikingDB のインデックス設定を調整してください。
- VikingDB を活用した多様な検索テストの実施
データをリアルタイムで VikingDB に書き込んだ後、ビジネスのシナリオに応じて異なる検索方式を選択できます:
- 検索方式:適用シーン / 詳細ドキュメント
- 多模態検索:画像からの類似画像検索、テキストによる画像検索、画像とテキストを混合した検索。画像・テキスト・混合入力をサポートします。 / https://docs.volcengine.com/docs/84313/1791135?lang=zh
- キーワード検索:完全一致や用語・番号類のクエリに最適。全文検索と意味検索を補完し合います。 / https://docs.volcengine.com/docs/84313/1791139
- 地理情報検索:地理位置および距離に基づいた検索フィルタリングをサポートします。 / https://docs.volcengine.com/docs/84313/1791133?lang=zh#filter%E7%BB%93%E6%9E%84
最後に
多模態 AI アプリケーションが本番環境で運用されるようになると、その価値はモデルの性能だけでなく、データの更新速度にも左右されます。画像、動画、音声、ドキュメントといった非構造化データにおいても、オブジェクトの内容やメタデータを Flink に接続できれば、「イベントトリガー型」「全量と増分を一体化した」「書き込み即検索可能」というアプローチで、企業向け AI アプリケーションのリアルタイムベクトル化パイプラインを構築することが可能です。
このパイプラインがもたらす核心的な変化は以下の通りです:
- イベント駆動によるスキャンの置換: ファイルアップロード後、数秒以内に処理がトリガーされます。データの可視化遅延は時間単位から数秒レベルに短縮されました。
- 単一の SQL ジョブで全ルートを貫通: 全量データ、増分データ、ベクトル化、そしてストレージを一つの SQL で完結させることで、アーキテクチャの複雑さが劇的に低下しました。
- ベクトル化能力は即戦力かつ自律可能: VikingDB の内蔵モデルを使えばゼロからエンジニアリングする必要もありませんし、Flink 2.2 AI SQL を用いて方舟(Ark)を呼び出すことで、モデルの自主的な運用も可能です。
- 書き込み後数秒で検索可能: 生成された Collection はそのまま、ナレッジベースの質問応答、推薦時の召回、そして多模態検索を支えることができます。
これはつまり、企業ナレッジベースの更新スピードが向上し、コンテンツ推薦が新しい素材を即座に感知できるようになり、トレーニング用サンプルもより迅速に蓄積されることを意味します。
火山引擎(Volcengine)の Flink と VikingDB を組み合わせることで、「多様なソース」「多様なモダリティ」「絶えず変化する」データをリアルタイムで検索可能かつサービス可能なベクトル資産へと変換。これにより、企業の AI アプリケーションは PoC(概念実証)から着実に本番環境へと移行していくことができます。
WeChat で開くにはこちらへ
原文を表示
原创 Viking 2026-08-04 19:12 北京
image
当企业 AI 应用从概念验证走向生产,一定遇到过这样的场景:商品图库每天新增几千张图片、企业知识库持续有新文档进入、训练数据平台需尽快感知新样本……这些内容的第一站,通常是对象存储。
但“文件已经上传”并不等于“内容已经能被 AI 使用”。从“存起来”到“用起来”,通常还要经过这样一套流程:
发现新增或变化的对象 → 读取对象或元信息 → 清洗与组装模型输入 → 生成 Embedding(向量化) → 写入向量数据库 → 更新检索索引
如果上述流程由多个定时任务和脚本拼接,生产阶段常常会遇到三个卡点:
数据更新不及时:新文件需要等待下一个扫描周期,检索内容可能滞后数小时甚至更久。
存量与增量难以衔接:全量扫描期间仍有新文件上传,切换增量消费时容易遗漏或重复。
模型和向量写入链路复杂:图片拉取、模型服务、GPU 资源、失败重试、向量写入与索引更新需要分别建设。
为了解决上述问题,本文将介绍如何用火山引擎 Flink 与 VikingDB 搭建一条“文件上传后秒级可检索”的实时多模态数据链路,并给出两套完整 Flink SQL 参考方案。
一、一条链路,收敛所有环节
Flink + VikingDB 联合方案将上述分散的环节收敛到一条持续运行的实时数据链路中:
这条链路的几个关键角色:
TOS(对象存储):承载图片、视频、文本、文档及业务元数据——你的多模态数据就存在这里。
Kafka(消息队列):承接 TOS 的 PUT、DELETE 等对象事件——文件变动时,Kafka 会收到一条消息。
流式计算 Flink 版:负责全增量接入、清洗、路由、模型调用与故障恢复——整条链路的“编排引擎”。
VikingDB(向量数据库):负责向量化、向量存储、索引和在线检索——向量数据的最终归宿,也是检索服务的后端。
方舟:为需要自定义模型的场景提供多模态 Embedding 能力——当内置模型不够用时,在这里接入自有模型。
二、三项关键能力,打通实时 AI 数据链路
1、TOS-CDC:把对象存储变成一张持续更新的表
TOS-CDC 是面向对象存储的 Flink SQL Source Connector。原本散落在对象存储里的文件变化,被连续地翻译成了一条可处理的数据流。作业启动后,它会:
记录全量扫描开始时间;
对指定的 TOS 存储桶执行全量扫描;
全量完成后,默认从全量扫描开始时间对应的 Kafka 位置开始消费对象事件;
通过 Checkpoint 保存全量扫描进度和 Kafka 消费位点。
通过从全量扫描时间开始衔接消费 Kafka 增量数据,能覆盖到全量扫描期间发生的对象变化。下游使用稳定主键 Upsert(即"有则更新,无则插入")后,即使有重复事件也最终能收敛到同一条记录。
注:TOS-CDC 接收到的是文件在哪里的信息,不接收文件内容本身。
2、VikingDB:让写入、向量化和检索形成闭环
对于标准图片和文本检索场景,可以在 VikingDB 表上声明字段语义和向量模型。例如将字段声明为 image,并配置 doubao-embedding-vision,由 VikingDB 自动完成图片读取、向量化和索引更新。
业务不需要额外维护模型服务、GPU 资源池和向量导入程序。Flink 负责持续写入变化的数据,VikingDB 将其沉淀为可检索的向量资产。
VikingDB Connector 支持根据 Flink Changelog 执行 Upsert 和 Delete。默认使用同步写入:
'async' = 'false'
对于要求写入后快速检索的业务,不建议直接启用 async=true。异步写入更偏向吞吐优先,会增加 Collection(数据集合)与 Index 的可见延迟。
3、 Flink 2.2 AI SQL:把模型调用变成 SQL 的一部分
部分业务需要在 Flink 侧自行完成 Embedding,典型场景包括:
将图片、标题、标签和 OCR 文本组装成多模态输入;
使用指定的方舟模型及版本;
统一控制模型调用的并行度、吞吐和成本;
让一份 Embedding 同时写入 VikingDB、Kafka、特征库或训练样本;
在不重写 Java/Python 作业的情况下切换模型。
流式计算 Flink 版 2.2 支持通过 CREATE MODEL 声明方舟模型,并使用 ML_PREDICT 在 SQL 中执行实时推理。模型接入、数据处理和向量写入由一条 SQL 作业统一编排。
注:当前流式计算 Flink 版 2.2 AI SQL 正在邀测中。如有需求可联系火山官网。
三、前期准备工作 CheckList
开始搭建链路前,需要准备以下资源:
准备项
作用
是否必需
TOS Bucket 与事件通知规则
保存多模态对象,并将对象变更事件投递至 Kafka
必需
Kafka 实例、Topic 与用户
承接 TOS 增量事件,供 TOS-CDC 持续消费
必需
流式计算 Flink 版
运行 TOS-CDC、数据处理和 VikingDB Sink
必需
TOS-CDC 邀测资格
使用全量扫描与增量事件一体化 Source
必需
Flink 2.2 邀测资格
使用 CREATE MODEL 与 ML_PREDICT 调用方舟
方案二必需
VikingDB 实例与 API Key
完成自动向量化或存储 Flink 生成的向量
必需
方舟推理接入点与 API Key
执行自定义多模态 Embedding
方案二必需
注:两种方案后续会展开介绍。
1、准备 TOS Bucket,并开启事件投递
TOS 事件通知能够在 Bucket 内对象发生变化时,将事件消息推送至 Kafka。事件消息包含 Bucket、对象 Key、事件类型和事件时间等信息,TOS-CDC 根据这些消息持续感知新增、覆盖和删除操作。
在 TOS 控制台为目标 Bucket 创建事件通知规则时,需要:
将推送目标设置为消息队列 Kafka 版;
至少订阅 tos:ObjectCreated:*;如果后续需要处理对象删除,再订阅 tos:ObjectRemoved:*;
按业务范围配置 Prefix、Suffix,使事件范围与 TOS-CDC 的 bucket 配置保持一致;
选择目标 Kafka 实例、Topic、Kafka 用户和授权角色。
注:详细操作请参见火山引擎对象存储文档:设置事件通知推送至 Kafka。(复制链接至浏览器:https://docs.volcengine.com/docs/6349/1817509?lang=zh)
2、准备 Kafka 实例、Topic 与访问授权
在消息队列 Kafka 版中创建实例、Topic 和访问用户。TOS 事件通知规则需要引用 Kafka 实例 ID、Topic 名称、用户和 IAM 角色。该角色需要绑定系统预设策略 KafkaAccessForTOS,用于授权 TOS 向 Kafka 投递事件。
同时需要确保:
Flink 资源池能够访问 Kafka Bootstrap Servers;
Kafka 用户与认证参数可以在 Flink 作业中使用;
Topic Retention 大于“全量扫描最大耗时 + rewind.offset”,避免全量扫描结束时需要回拨的事件已经过期。
3、开通流式计算 Flink 版
开通火山引擎流式计算 Flink 版,创建项目和运行作业所需的资源池,并打通到 Kafka、VikingDB 及方舟服务的网络。
注:若需 TOS-CDC 邀测资格 和 Flink 2.2 版本邀测资格,可联系火山官网。
4、准备 VikingDB 与方舟资源
开通 VikingDB,准备数据面地址和 API Key。方案一需要确认目标 Collection 使用的自动向量化模型、版本和维度;方案二需要创建方舟推理接入点,准备模型名称、输出维度和 API Key。
注:所有访问凭证均建议通过流式计算 Flink 版的加密变量或运行环境变量注入,不要直接写入 SQL。
三、两套 SQL 方案:选你需要的那条路
两套方案共用同一条 TOS-CDC 数据接入链路,区别在于“谁来完成向量化操作”:
方案一:VikingDB 自动向量化。 面向标准图文检索,架构最简单——把数据交给 VikingDB,它来搞定 Embedding。
方案二:Flink AI SQL 调用方舟。 面向图文融合、模型自主和向量多下游复用——用户自主控制模型、输入和输出。
1、创建 TOS-CDC 源表
CREATE TABLE tos_object_events (
object_key STRING NOT NULL,
object_url STRING,
bucket_name STRING,
file_name STRING,
object_etag STRING,
object_size BIGINT,
mtime TIMESTAMP_LTZ(3),
event_time TIMESTAMP_LTZ(3),
record_origin STRING,
PRIMARY KEY (object_key) NOT ENFORCED
) WITH (
'connector' = 'tos-cdc',
'path' = 'tos://my-bucket/images01/, tos://my-bucket/images02/',
'properties.bootstrap.servers' = 'kafka.example:9092',
'properties.group.id' = 'ingest-cg',
'topic' = 'object-events',
'scan.startup.mode' = 'initial'
);
生产环境需要根据 Kafka 实例补充认证和网络参数,并注意:
Kafka Topic 的消息保留时间应能覆盖全量阶段扫描到 Kafka 切换所需时间以及配置的 rewind.offset。
增量回拨参数 rewind.offset 默认为0,可按需设置,并建议保留 rewind.retention-miss-policy=fail,避免回拨位置过期后静默漏数。
构造输出数据:
CREATE TEMPORARY VIEW image_put_events AS
SELECT
object_key AS id,
object_url AS image_uri,
object_etag,
COALESCE(event_time, mtime) AS update_time
FROM tos_object_events;
这里使用 object_key 作为主键,因 TOS-CDC 与外部 Sink 均采用 At-Least-Once 语义(至少投递一次,可能重复),故障恢复或全增量衔接期间可能重放记录。稳定主键可以使重复 PUT 在 VikingDB 中执行 Upsert,最终收敛到同一条数据。
2、方案一:VikingDB 自动向量化
首先创建 VikingDB Catalog:
CREATE CATALOG viking WITH (
'type' = 'vikingdb',
'control-plane.host' = 'open.volcengineapi.com',
'region' = 'cn-beijing',
'project-name' = 'default',
'access-key' = '${secret_values.volc-ak}',
'secret-key' = '${secret_values.volc-sk}',
'data-plane.host' = '<vikingdb-data-plane-host>',
'api-key' = '${secret_values.vikingdb-api-key}'
);
然后创建启用自动图片向量化的 Collection:
CREATE TABLE IF NOT EXISTS viking.default.realtime_image_assets (
id STRING,
image_uri STRING,
object_etag STRING,
update_time TIMESTAMP_LTZ(3),
PRIMARY KEY (id) NOT ENFORCED
) WITH (
'vikingdb.field.image_uri.type' = 'image',
'vikingdb.vectorize.dense.model-name' = 'doubao-embedding-vision',
'vikingdb.vectorize.dense.model-version' = '<model-version>',
'vikingdb.vectorize.dense.dim' = '2048',
'vikingdb.vectorize.dense.image-field' = 'image_uri'
);
最后写入 VikingDB:
INSERT INTO viking.default.realtime_image_assets
SELECT id, image_uri, object_etag, update_time
FROM image_put_events;
这条链路能够覆盖:
作业首次启动时导入指定存储桶下的存量对象;
新对象上传后实时写入;
同一路径对象被覆盖后按稳定主键更新;
作业失败后从 Checkpoint 恢复,并通过 Upsert 抵御事件重放。
3、方案二:Flink AI SQL 调用方舟 Embedding
当业务需要自定义多模态输入或指定模型时,可以复用同一张 TOS-CDC 源表。
首先在 SQL 中声明方舟模型:
CREATE MODEL ark_multimodal_embedding
INPUT (payload STRING)
OUTPUT (embedding ARRAY<FLOAT>)
WITH (
'provider' = 'ark',
'endpoint' = 'https://ark.cn-beijing.volces.com/api/v3/embeddings/multimodal',
'api-key' = '${secret_values.ark-api-key}',
'model' = 'doubao-embedding-vision-251215',
'model.dimensions' = '2048'
);
将对象信息组装为方舟多模态输入:
CREATE TEMPORARY VIEW multimodal_payload AS
SELECT
id,
image_uri,
update_time,
CAST(
JSON_ARRAY(
JSON_OBJECT(
'type' VALUE 'image_url',
'image_url' VALUE JSON_OBJECT(
'url' VALUE CONCAT(
'https://<bucket-domain>/',
object_key
)
)
)
) AS STRING
) AS payload
FROM (
SELECT
object_key AS id,
object_url AS image_uri,
COALESCE(event_time, mtime) AS update_time,
object_key
FROM tos_object_events
) AS source_events;
示例使用 HTTPS 图片地址作为模型输入。生产环境应确保方舟服务能够安全访问该地址;私有 Bucket 可以使用受控的临时签名 URL 或企业内部授权链路,不建议为模型调用将整个 Bucket 配置为公开读。
创建保存显式向量的 VikingDB Collection:
CREATE TABLE IF NOT EXISTS viking.default.realtime_multimodal_assets (
id STRING,
image_uri STRING,
update_time TIMESTAMP_LTZ(3),
mixed_embedding ARRAY<FLOAT>,
PRIMARY KEY (id) NOT ENFORCED
) WITH (
'vikingdb.field.mixed_embedding.type' = 'vector',
'vikingdb.field.mixed_embedding.dim' = '2048'
);
调用模型并写入 VikingDB:
INSERT INTO viking.default.realtime_multimodal_assets
SELECT
id,
image_uri,
update_time,
embedding AS mixed_embedding
FROM ML_PREDICT(
TABLE multimodal_payload,
MODEL ark_multimodal_embedding,
DESCRIPTOR(payload)
);
如果 Embedding 还要用于实时特征、训练样本或消息订阅,可以通过 EXECUTE STATEMENT SET 增加多个 Sink,让下游共享同一次模型计算结果。
五、如何验证这条链路
1、确认 Flink 任务进入运行状态
通过 Flink UI 检查,确认存量的图片、视频文件已经导入 VikingDB。确保数据量和 TOS 能够对齐。
2、验证向量与搜索结果
在 VikingDB 控制台的"数据集 → 数据预览"中,按照 TOS 的路径进行查询,确认数据已经写入数据集。
检查 VikingDB Collection 中的字段类型和向量维度,并使用一张相似图片或一段相关文本发起检索,确认能够召回刚上传的对象。如下图所示,输入“小松鼠”可以召回相关相似的照片。
3、测量端到端时延
分别记录:
TOS 对象上传时间;
Kafka 事件时间;
Flink 处理时间;
VikingDB 写入可见时间;
首次能够检索到该对象的时间。
以真实数据规模和并发条件评估 P50、P95 延迟,再调整 Flink 并行度、模型吞吐、Sink Flush Interval 和 VikingDB 索引配置。
4、使用 VikingDB 做多样化检索测试
数据实时写入 VikingDB 后,可根据业务场景选择不同检索方式:
检索方式
适用场景
详细文档
多模态检索
以图搜图、以文搜图、图文混合召回,支持图片/文本/混合输入。
https://docs.volcengine.com/docs/84313/1791135?lang=zh
关键词检索
精确匹配、术语/编号类查询,实现全文检索与语义检索互补。
https://docs.volcengine.com/docs/84313/1791139
地理信息检索
支持按地理位置和距离做检索过滤。
https://docs.volcengine.com/docs/84313/1791133?lang=zh#filter%E7%BB%93%E6%9E%84
写在最后
多模态 AI 应用进入生产阶段后,价值不只来自模型效果,也来自数据更新速度。对于图片、视频、音频、文档等非结构化数据,只要能把对象内容或元信息接入 Flink,就可以沿用“事件触发、全增量一体、写入即可检索”的方式,构建面向企业 AI 应用的实时向量化链路。
这条链路带来的核心改变:
事件驱动替代定时扫描:文件上传后秒级触发处理,数据可见延迟从小时级降至秒级。
一条 SQL 作业打通全链路:全量 + 增量 + 向量化 + 存储,架构复杂度大幅下降。
向量化能力开箱即用,又可自主可控:既能用 VikingDB 内置模型零工程落地,也能用 Flink 2.2 AI SQL 调用方舟实现模型自主。
写入秒级可见、可搜索:产出的 Collection 直接支撑知识库问答、推荐召回与多模态检索。
这意味着企业知识库能更快更新、内容推荐能更快感知新素材、训练样本也能更及时沉淀。
火山引擎 Flink + VikingDB,把“多源、多模态、持续变化”的数据实时转化为可检索、可服务的向量资产,助力企业 AI 应用从 PoC 稳步迈向生产。
跳转微信打开
今日のまとめ
AIデイリーブリーフで今日の重要ニュースをまとめ読み