All Products
Search
Document Center

Realtime Compute for Apache Flink:Memulai analitik real-time berbasis AI

Last Updated:Jul 18, 2026

Topik ini menjelaskan cara menggunakan model bawaan Flink AI untuk analitik data real-time.

Informasi latar belakang

阿里云百炼大模型服务平台 adalah platform satu atap untuk pengembang AI. Platform ini terintegrasi secara mendalam dengan Realtime Compute for Apache Flink. Anda dapat dengan cepat memanggil kemampuan LLM dan menggabungkannya dengan pipa data real-time Flink untuk memperpendek jalur dari ingesti data hingga pengambilan keputusan cerdas. Dua jenis model inti didukung:

  • Model chat/completions: LLM generatif untuk pemahaman teks, banyak digunakan untuk analisis sentimen, pengenalan maksud, dan tanya jawab (Q&A).

    • Analisis sentimen: Klasifikasikan komentar media sosial secara real-time untuk mengidentifikasi sentimen pengguna (positif, negatif, netral).

    • Layanan pelanggan cerdas: Berikan interaksi bahasa alami melalui dialog generatif.

    • Moderasi konten: Deteksi secara otomatis konten sensitif atau tidak sesuai dalam teks.

  • Model embedding: Mengonversi teks menjadi representasi vektor berdimensi tinggi untuk pencarian semantik, sistem rekomendasi, dan pembangunan graf pengetahuan.

    • Pencarian semantik: Vektorisasi deskripsi produk atau kueri pengguna untuk pencarian relevansi berbasis semantik.

    • Sistem rekomendasi: Manfaatkan vektorisasi teks untuk menemukan korelasi laten antara minat pengguna dan fitur produk.

    • Graf pengetahuan: Konversi teks tak terstruktur ke bentuk vektor untuk ekstraksi pengetahuan dan pemodelan relasi.

Prasyarat

Batasan

Hanya didukung oleh mesin komputasi real-time VVR 11.1 dan versi yang lebih tinggi.

Langkah 1: Daftarkan model

Daftarkan model dengan mengikuti panduan di Configure a model.

Model chat/completions

SQL berikut mendaftarkan model Qwen-turbo:

CREATE MODEL ai_analyze_sentiment
INPUT (`input` STRING)
OUTPUT (`content` STRING)
WITH (
    'provider'='openai-compat',
    'task' = 'chat/completions',
    'model'='qwen-turbo',                                                               -- Qwen-turbo model
    'system-prompt' = 'Classify the text below into one of the following labels: [positive, negative, neutral, mixed]. Output only the label.'
);

Model embedding

SQL berikut mendaftarkan model text-embedding-v3:

CREATE MODEL embedding_model
INPUT (`input` STRING)
OUTPUT (`embeddings` ARRAY<FLOAT>)
WITH (
    'provider'='openai-compat',
    'task' = 'embeddings',
    'model'='text-embedding-v3'                                                  -- text-embedding-v3 model
);

Langkah 2: Buat pekerjaan

Buat draf pekerjaan streaming SQL. Untuk informasi selengkapnya, lihat Flink SQL jobs.

Langkah 3: Tulis pekerjaan SQL untuk analitik AI

Model chat/completions

Gunakan ML_PREDICT untuk memanggil model ai_analyze_sentiment yang telah didaftarkan guna melakukan analisis sentimen pada komentar film.

Penting

Throughput operator ML_PREDICT tunduk pada pembatasan kecepatan. Saat batas tercapai, tekanan balik terjadi dengan operator ML_PREDICT sebagai bottleneck. Pembatasan kecepatan yang parah dapat menyebabkan timeout operator dan restart pekerjaan. Untuk detail pembatasan kecepatan, lihat 限流. Hubungi manajer bisnis Anda untuk meminta peningkatan batas.

Salin SQL berikut ke editor SQL.

-- Create a temporary result table
CREATE TEMPORARY TABLE print_sink(
  id BIGINT,
  movie_name VARCHAR, 
  predict_label VARCHAR, 
  actual_label VARCHAR
) WITH (
  'connector' = 'print',   -- print connector
  'logger' = 'true'        -- display results in console
);
-- Create a temporary view with test data
-- | id | movie_name | comment   | actual_label |
-- | 1  | 好东西     | 最爱小孩子猜声音那段,算得上看过的电影里相当浪漫的叙事了。很温和也很有爱。| POSITIVE |
-- | 2  | 水饺皇后   | 乏善可陈  | NEGATIVE |
CREATE TEMPORARY VIEW movie_comment(id, movie_name, user_comment, actual_label)
AS VALUES (1, '好东西', '最爱小孩子猜声音那段,算得上看过的电影里相当浪漫的叙事了。很温和也很有爱。', 'positive'), (2, '水饺皇后', '乏善可陈', 'negative');
INSERT INTO print_sink
SELECT id, movie_name, content as predict_label, actual_label 
FROM ML_PREDICT(
  TABLE movie_comment, 
  MODEL ai_analyze_sentiment,  -- The registered Qwen-turbo model
  DESCRIPTOR(user_comment));   

Model embedding

Gunakan ML_PREDICT untuk memanggil model embedding_model yang telah didaftarkan guna melakukan vektorisasi komentar film dan menuliskan hasilnya ke Milvus (pratinjau publik).

Penting

Throughput operator ML_PREDICT tunduk pada pembatasan kecepatan. Saat batas tercapai, tekanan balik terjadi dengan operator ML_PREDICT sebagai bottleneck. Pembatasan kecepatan yang parah dapat menyebabkan timeout operator dan restart pekerjaan. Untuk detail pembatasan kecepatan, lihat 限流. Hubungi manajer bisnis Anda untuk meminta peningkatan batas.

Salin SQL berikut ke editor SQL.

-- Create a temporary Milvus sink table
CREATE TEMPORARY TABLE milvus_sink
(
    id STRING,
    movie_name STRING,
    user_comment STRING,
    embeddings ARRAY<FLOAT>,
    PRIMARY KEY (id) NOT ENFORCED
)
WITH (
    'connector' = 'milvus',
    'endpoint' = '<YOUR-ENDPOINT>',
    'port' = '<YOUR-PORT>',
    'userName' = '<YOUR-USERNAME>',
    'password' = '<YOUR-PASSWORD>',
    'databaseName' = 'default',
    'collectionName' = 'movie-comment-embeddings'
);
-- Create a temporary view with test data
-- | id | movie_name | comment | actual_label|
-- | 1 | 好东西 |最爱小孩子猜声音那段,算得上看过的电影里相当浪漫的叙事了。很温和也很有爱。| POSITIVE |
-- | 2 | 水饺皇后 | 乏善可陈 | NEGATIVE |
CREATE TEMPORARY VIEW movie_comment(id, movie_name,  user_comment)
AS VALUES ('1', '好东西', '最爱小孩子猜声音那段,算得上看过的电影里相当浪漫的叙事了。很温和也很有爱。'), ('2', '水饺皇后', '乏善可陈');
INSERT INTO
    milvus_sink
SELECT
    id,
    movie_name,
    user_comment,
    embeddings
FROM
    ML_PREDICT (
        TABLE movie_comment,
        MODEL embedding_model,  -- The registered text-embedding-v3 model                   
        DESCRIPTOR (user_comment)
    );

Langkah 4: Deploy dan mulai pekerjaan

Deploy dan mulai pekerjaan. Untuk informasi selengkapnya, lihat Flink SQL jobs.

Langkah 5: Lihat hasil analisis

Model chat/completions

  1. Verifikasi bahwa status pekerjaan adalah Completed.

  2. Pada halaman O&M Center > Job O&M, klik nama pekerjaan target.

  3. Pada tab Job Logs, pilih Task Managers dan pilih TaskManager saat ini.

  4. Klik Logs dan cari entri PrintSinkOutputWriter.

    Prediksi model predict_label sesuai dengan actual_label.

    Pilih tab Task Managers dan pilih Logs. Output menampilkan hasil seperti +I[1, 好东西, positive, positive] dan +I[2, 水饺皇后, negative, negative], yang mengonfirmasi bahwa predict_label sesuai dengan actual_label.

Model embedding

Hasil dituliskan ke koleksi Milvus target. Masuk ke Milvus untuk melihatnya.

  1. Mengakses halaman Attu.

  2. Buka koleksi target dan lihat data yang telah disinkronkan pada tab Data.

    Setelah sinkronisasi, tab Data menampilkan 1.000 catatan yang telah disinkronkan dengan bidang termasuk id, vector, dan timestamp.