All Products
Search
Document Center

Realtime Compute for Apache Flink:ApsaraMQ for RocketMQ

Last Updated:Apr 25, 2026

Topik ini memperkenalkan konektor ApsaraMQ for RocketMQ.

Penting

Instans ApsaraMQ for RocketMQ Edisi Standar 4.x memiliki batas atas elastis sebesar 5.000 panggilan API per detik. Melebihi batas ini saat menghubungkan instans ke Realtime Compute for Apache Flink akan memicu throttling, yang dapat mengganggu stabilitas pekerjaan Flink Anda. Oleh karena itu, jika Anda menggunakan atau berencana menggunakan instans RocketMQ Edisi Standar untuk integrasi dengan Flink, evaluasi secara cermat dampak potensialnya. Jika memungkinkan, pertimbangkan middleware messaging alternatif, seperti Kafka, Simple Log Service (SLS), atau DataHub. Jika Anda harus menggunakan instans ApsaraMQ for RocketMQ Edisi Standar 4.x untuk volume pesan tinggi, submit a ticket untuk meminta kuota throttling yang lebih tinggi.

Informasi latar belakang

ApsaraMQ for RocketMQ adalah middleware messaging terdistribusi yang dikembangkan oleh Alibaba Cloud berdasarkan Apache RocketMQ. Layanan ini menyediakan latensi rendah, konkurensi tinggi, ketersediaan tinggi, dan keandalan tinggi. ApsaraMQ for RocketMQ menawarkan penguraian keterkaitan asinkron dan peak shaving untuk sistem aplikasi terdistribusi, serta fitur-fitur untuk aplikasi Internet seperti akumulasi pesan dalam jumlah besar, throughput tinggi, dan retry yang andal.

Tabel berikut menjelaskan konektor ApsaraMQ for RocketMQ.

Item

Deskripsi

Jenis yang didukung

tabel sumber dan tabel sink

Mode operasi

hanya mode streaming

Format data

format CSV dan biner

Metrik khusus konektor

Metrik

  • Tabel sumber

    • numRecordsIn

    • numRecordsInPerSecond

    • numBytesIn

    • numBytesInPerSecond

    • currentEmitEventTimeLag

    • currentFetchEventTimeLag

    • sourceIdleTime

  • Tabel sink

    • numRecordsOut

    • numRecordsOutPerSecond

    • numBytesOut

    • numBytesOutPerSecond

    • currentSendTime

Catatan

Lihat Metrik untuk detail selengkapnya.

Jenis API

DataStream API (hanya untuk RocketMQ 4.x) dan SQL API

Dukungan pembaruan atau penghapusan data pada tabel sink

Hanya penyisipan data ke tabel sink yang didukung. Pembaruan dan penghapusan tidak didukung.

Fitur

Tabel sumber dan sink ApsaraMQ for RocketMQ mendukung bidang metadata berikut.

  • Bidang untuk tabel sumber

    Bidang

    Tipe

    Deskripsi

    topic

    VARCHAR METADATA VIRTUAL

    Topik pesan.

    queue-id

    INT METADATA VIRTUAL

    ID antrian.

    queue-offset

    BIGINT METADATA VIRTUAL

    Offset konsumsi.

    msg-id

    VARCHAR METADATA VIRTUAL

    ID pesan.

    store-timestamp

    TIMESTAMP(3) METADATA VIRTUAL

    Timestamp penyimpanan pesan.

    born-timestamp

    TIMESTAMP(3) METADATA VIRTUAL

    Timestamp pembuatan pesan.

    keys

    VARCHAR METADATA VIRTUAL

    Kunci pesan.

    tags

    VARCHAR METADATA VIRTUAL

    Tag pesan.

  • Bidang untuk tabel sink

    Bidang

    Tipe

    Deskripsi

    keys

    VARCHAR METADATA

    Kunci pesan.

    tags

    VARCHAR METADATA

    Tag pesan.

Prasyarat

Anda telah membuat resource Message Queue for Apache RocketMQ. Untuk petunjuknya, lihat Buat resource.

Batasan

  • ApsaraMQ for RocketMQ 5.x memerlukan mesin komputasi waktu nyata Flink VVR 8.0.3 atau versi yang lebih baru.

  • Konektor ApsaraMQ for RocketMQ menggunakan pull consumer, yang mendistribusikan beban kerja di seluruh subtask.

Sintaks

CREATE TABLE mq_source(
  x varchar,
  y varchar,
  z varchar
) WITH (
  'connector' = 'mq5',
  'topic' = '<yourTopicName>',
  'endpoint' = '<yourEndpoint>',
  'consumerGroup' = '<yourConsumerGroup>'
);

Parameter WITH

Umum

Parameter

Deskripsi

Tipe

Wajib

Default

Keterangan

connector

Jenis konektor.

String

Ya

None

  • Untuk RocketMQ 4.x, nilainya adalah mq.

  • Untuk RocketMQ 5.x, nilainya adalah mq5.

endPoint

Titik akhir layanan.

String

Ya

None

ApsaraMQ for RocketMQ menyediakan dua jenis titik akhir:

  • Titik akhir untuk layanan MQ pada jaringan internal (jaringan klasik atau VPC): Pada halaman detail instans target di Konsol MQ, pilih Endpoints > TCP Protocol Client Endpoints > Internal Network Access untuk mendapatkan titik akhir yang sesuai.

  • Titik akhir layanan MQ publik: Pada halaman detail instans target di Konsol MQ, pilih Endpoint > TCP Protocol > Client Endpoint > Public Access untuk mendapatkan titik akhir yang sesuai.

Penting

Kami menyarankan Anda menggunakan titik akhir VPC. Koneksi jaringan publik mungkin tidak stabil karena perubahan dinamis pada kebijakan keamanan jaringan Alibaba Cloud.

  • Jaringan internal tidak mendukung akses lintas wilayah. Misalnya, jika layanan Realtime Compute for Apache Flink Anda berada di wilayah China (Hangzhou) dan instans Alibaba Cloud Message Queue for Apache RocketMQ Anda berada di wilayah China (Shanghai), koneksi akan gagal.

  • Untuk menghubungkan melalui jaringan publik, Anda harus mengaktifkan akses publik untuk instans tersebut. Untuk informasi selengkapnya, lihat Koneksi jaringan.

topic

Nama topik.

String

Ya

None

None

accessId

  • Untuk RocketMQ 4.x: ID AccessKey Akun Alibaba Cloud Anda.

  • Untuk RocketMQ 5.x:

    Username instans RocketMQ.

String

  • Untuk RocketMQ 4.x: Ya

  • Untuk RocketMQ 5.x: Tidak

None

Penting

Untuk menghindari eksposur pasangan AccessKey Anda, kami menyarankan menggunakan variabel proyek untuk menentukan ID AccessKey dan Secret AccessKey.

  • RocketMQ 5.x:

    • Anda menggunakan titik akhir publik.

    • Anda menggunakan titik akhir VPC dan akses bebas otentikasi melalui jaringan internal dinonaktifkan.

    • Pengaturan ini tidak diperlukan jika Anda menggunakan titik akhir VPC dan akses bebas otentikasi jaringan internal diaktifkan.

accessKey

  • Untuk RocketMQ 4.x: Secret AccessKey Akun Alibaba Cloud Anda.

  • Untuk RocketMQ 5.x: Password instans.

String

  • Untuk RocketMQ 4.x: Ya

  • Untuk RocketMQ 5.x: Tidak

None

tag

Tag pesan yang akan berlangganan atau ditulis.

String

Tidak

None

  • Saat RocketMQ digunakan sebagai sumber, Anda dapat membaca pesan dengan satu tag.

  • Saat RocketMQ digunakan sebagai sink, Anda dapat menentukan beberapa tag, dipisahkan dengan koma (,).

Catatan

Saat digunakan sebagai sink, parameter ini hanya didukung untuk RocketMQ 4.x. Untuk RocketMQ 5.x, tentukan tag pesan di bidang metadata sink.

encoding

Format encoding.

String

Tidak

UTF-8

None

instanceID

ID instans Alibaba Cloud Message Queue for Apache RocketMQ.

String

Tidak

None

  • Jika instans tidak memiliki namespace khusus, jangan konfigurasikan parameter instanceID.

  • Jika instans memiliki namespace khusus, parameter instanceID wajib diisi.

Catatan

Parameter ini hanya didukung untuk RocketMQ 4.x.

Khusus sumber

Parameter

Deskripsi

Tipe

Wajib

Default

Keterangan

consumerGroup

Nama kelompok konsumen.

String

Ya

None

None

pullIntervalMs

Interval polling dalam milidetik untuk sumber saat tidak ada data yang tersedia.

Int

Ya

None

Unit: milidetik.

Mekanisme throttling tidak tersedia. Anda tidak dapat mengatur laju pembacaan data dari RocketMQ.

Catatan

Parameter ini hanya didukung untuk RocketMQ 4.x.

timeZone

Zona waktu.

String

Tidak

None

Contoh: Asia/Shanghai.

startTimeMs

Waktu mulai untuk konsumsi data.

Long

Tidak

None

Timestamp dalam milidetik.

startMessageOffset

Offset pesan tempat konsumsi dimulai.

Int

Tidak

None

Jika parameter ini ditentukan, pemuatan data dimulai dari offset yang ditentukan oleh startMessageOffset, yang memiliki prioritas lebih tinggi.

lineDelimiter

Pembatas baris yang digunakan untuk mengurai catatan.

String

Tidak

\n

None

fieldDelimiter

Pembatas bidang.

String

Tidak

\u0001

Pembatas bervariasi berdasarkan mode terminal:

  • Dalam mode read-only (default), pembatasnya adalah \u0001. Pembatas ini tidak terlihat dalam mode ini.

  • Dalam mode edit, pembatasnya adalah ^A.

lengthCheck

Kebijakan untuk memeriksa jumlah bidang dalam setiap catatan.

String

Tidak

NONE

Nilai yang valid:

  • NONE: Nilai default.

    • Jika catatan memiliki lebih banyak bidang daripada yang didefinisikan skema, bidang tambahan di sebelah kanan dipotong.

    • Jika catatan memiliki lebih sedikit bidang daripada yang didefinisikan skema, catatan tersebut dilewati.

  • SKIP: Melewatkan catatan apa pun di mana jumlah bidang tidak sesuai dengan skema.

  • EXCEPTION: Memicu exception jika jumlah bidang tidak sesuai dengan skema.

  • PAD: Mengisi bidang dari kiri ke kanan.

    • Jika catatan memiliki lebih banyak bidang daripada yang didefinisikan skema, bidang tambahan di sebelah kanan dipotong.

    • Jika catatan memiliki lebih sedikit bidang daripada yang didefinisikan skema, bidang yang hilang di sebelah kanan diisi dengan nilai null.

columnErrorDebug

Menentukan apakah akan mengaktifkan mode debug untuk error parsing kolom.

Boolean

Tidak

false

Jika diatur ke true, log detail untuk exception parsing akan dicetak.

pullBatchSize

Jumlah maksimum pesan yang ditarik dalam satu batch.

Int

Tidak

64

Didukung di VVR 8.0.7 dan versi yang lebih baru.

Khusus sink

Parameter

Deskripsi

Tipe

Wajib

Default

Keterangan

producerGroup

Nama kelompok produsen.

String

Ya

None

None

retryTimes

Jumlah kali percobaan ulang untuk operasi penulisan yang gagal.

Int

Tidak

10

None

sleepTimeMs

Interval antar percobaan ulang, dalam milidetik.

Long

Tidak

5000

None

partitionField

Nama bidang yang digunakan sebagai kunci partisi.

String

Tidak

None

Parameter ini wajib jika parameter mode diatur ke partition.

Catatan

Didukung di VVR 8.0.5 dan versi yang lebih baru.

deliveryTimestampMode

Mode pengiriman untuk pesan tertunda. Parameter ini bekerja bersama parameter deliveryTimestampValue untuk menentukan kapan pesan tertunda dikirimkan.

String

Tidak

None

Nilai yang valid:

  • fixed: mode timestamp tetap.

  • relative: mode penundaan relatif.

  • field: mode berbasis bidang.

Catatan

Didukung di VVR 11.1 dan versi yang lebih baru.

deliveryTimestampType

Jenis referensi waktu untuk pesan tertunda.

String

Tidak

processing_time

Nilai yang valid:

  • event_time: event time.

  • processing_time: processing time.

Catatan

Didukung di VVR 11.1 dan versi yang lebih baru.

deliveryTimestampValue

Waktu pengiriman pesan tertunda.

Long

Tidak

None

Makna parameter ini bergantung pada nilai deliveryTimestampMode:

  • deliveryTimestampMode=fixed: Pesan ditunda hingga mencapai cap waktu yang ditentukan dalam milidetik. Jika waktu saat ini telah melewati cap waktu tersebut, pesan akan dikirimkan segera.

  • deliveryTimestampMode=relative: Durasi penundaan, dalam milidetik, relatif terhadap referensi waktu yang ditentukan oleh deliveryTimestampType.

  • deliveryTimestampMode=field: Parameter ini diabaikan. Waktu pengiriman ditentukan oleh nilai bidang yang ditentukan oleh deliveryTimestampField.

Catatan

Didukung di VVR 11.1 dan versi yang lebih baru.

deliveryTimestampField

Menentukan bidang yang digunakan sebagai waktu pengiriman untuk pesan tertunda. Tipe data harus BIGINT.

String

Tidak

None

Berlaku saat deliveryTimestampMode bernilai field.

Catatan

Didukung di VVR 11.1 dan versi yang lebih baru.

Pemetaan tipe

Tipe Flink

Tipe RocketMQ

BOOLEAN

STRING

VARBINARY

VARCHAR

TINYINT

INTEGER

BIGINT

FLOAT

DOUBLE

DECIMAL

Contoh

Contoh tabel sumber

  • Format CSV

    Asumsikan sebuah pesan berisi catatan data berikut dalam format CSV.

    1,name,male 
    2,name,female
    Catatan

    Pesan Message Queue for Apache RocketMQ dapat berisi nol atau lebih catatan data, dipisahkan oleh \n.

    Gunakan DDL berikut dalam pekerjaan Flink Anda untuk mendeklarasikan tabel sumber Message Queue for Apache RocketMQ.

    • RocketMQ 5.x

    CREATE TABLE mq_source(
      id varchar,
      name varchar,
      gender varchar,
      topic varchar metadata virtual
    ) WITH (
      'connector' = 'mq5',
      'topic' = 'mq-test',
      'endpoint' = '<yourEndpoint>',
      'consumerGroup' = 'mq-group',
      'fieldDelimiter' = ','
    );
    • RocketMQ 4.x

    CREATE TABLE mq_source(
      id varchar,
      name varchar,
      gender varchar,
      topic varchar metadata virtual
    ) WITH (
      'connector' = 'mq',
      'topic' = 'mq-test',
      'endpoint' = '<yourEndpoint>',
      'pullIntervalMs' = '1000',
      'accessId' = '${secret_values.ak_id}',
      'accessKey' = '${secret_values.ak_secret}',
      'consumerGroup' = 'mq-group',
      'fieldDelimiter' = ','
    );
  • Format biner

    • RocketMQ 5.x

      CREATE TEMPORARY TABLE source_table (
        mess varbinary
      ) WITH (
        'connector' = 'mq5',
        'endpoint' = '<yourEndpoint>',
        'topic' = 'mq-test',
        'consumerGroup' = 'mq-group'
      );
      
      CREATE TEMPORARY TABLE out_table (
        commodity varchar
      ) WITH (
        'connector' = 'print'
      );
      
      INSERT INTO out_table
      select 
        cast(mess as varchar)
      FROM source_table;
    • RocketMQ 4.x

      CREATE TEMPORARY TABLE source_table (
        mess varbinary
      ) WITH (
        'connector' = 'mq',
        'endpoint' = '<yourEndpoint>',
        'pullIntervalMs' = '500',
        'accessId' = '${secret_values.ak_id}',
        'accessKey' = '${secret_values.ak_secret}',
        'topic' = 'mq-test',
        'consumerGroup' = 'mq-group'
      );
      
      CREATE TEMPORARY TABLE out_table (
        commodity varchar
      ) WITH (
        'connector' = 'print'
      );
      
      INSERT INTO out_table
      select 
        cast(mess as varchar)
      FROM source_table;

Contoh tabel sink

  • Buat tabel sink

    • RocketMQ 5.x

      CREATE TABLE mq_sink (
        id INTEGER,
        len BIGINT,
        content VARCHAR
      ) WITH (
        'connector'='mq5',
        'endpoint'='<yourEndpoint>',
        'topic'='<yourTopicName>',
        'producerGroup'='<yourGroupName>'
      );
    • RocketMQ 4.x

      CREATE TABLE mq_sink (
        id INTEGER,
        len BIGINT,
        content VARCHAR
      ) WITH (
        'connector'='mq',
        'endpoint'='<yourEndpoint>',
        'accessId'='${secret_values.ak_id}',
        'accessKey'='${secret_values.ak_secret}',
        'topic'='<yourTopicName>',
        'producerGroup'='<yourGroupName>'
      );
      Catatan

      Untuk pesan RocketMQ dalam format biner, DDL harus mendefinisikan satu bidang dengan tipe data VARBINARY.

  • Buat tabel sink yang memetakan bidang metadata keys dan tags ke kunci dan tag pesan

    • RocketMQ 5.x

      CREATE TABLE mq_sink (
        id INTEGER,
        len BIGINT,
        content VARCHAR,
        keys VARCHAR METADATA,
        tags VARCHAR METADATA
      ) WITH (
        'connector'='mq5',
        'endpoint'='<yourEndpoint>',
        'topic'='<yourTopicName>',
        'producerGroup'='<yourGroupName>'
      );
    • RocketMQ 4.x

      CREATE TABLE mq_sink (
        id INTEGER,
        len BIGINT,
        content VARCHAR,
        keys VARCHAR METADATA,
        tags VARCHAR METADATA
      ) WITH (
        'connector'='mq',
        'endpoint'='<yourEndpoint>',
        'accessId'='${secret_values.ak_id}',
        'accessKey'='${secret_values.ak_secret}',
        'topic'='<yourTopicName>',
        'producerGroup'='<yourGroupName>'
      );

DataStream API

Penting

Untuk membaca dan menulis data menggunakan DataStream API, gunakan konektor DataStream yang sesuai untuk terhubung ke Realtime Compute for Apache Flink. Untuk informasi selengkapnya tentang konfigurasi konektor DataStream, lihat Integrasikan konektor DataStream.

VVR menyediakan MetaQSource untuk membaca dari RocketMQ dan MetaQOutputFormat, implementasi dari OutputFormat, untuk menulis ke RocketMQ. Contoh berikut menunjukkan cara membaca dari dan menulis ke RocketMQ:

RocketMQ 5.x

Catatan

Pada ApsaraMQ for RocketMQ 5.x, pasangan kunci akses merepresentasikan username dan password untuk instans. Anda tidak perlu mengonfigurasi pasangan ini jika mengakses instans melalui jaringan internal dan otentikasi ACL dinonaktifkan.

import com.alibaba.ververica.connectors.common.sink.OutputFormatSinkFunction;
import com.alibaba.ververica.connectors.mq5.shaded.org.apache.rocketmq.common.message.MessageExt;
import com.alibaba.ververica.connectors.mq5.sink.RocketMQOutputFormat;
import com.alibaba.ververica.connectors.mq5.source.RocketMQSource;
import com.alibaba.ververica.connectors.mq5.source.reader.deserializer.RocketMQRecordDeserializationSchema;
import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.api.common.functions.MapFunction;
import org.apache.flink.api.common.typeinfo.TypeInformation;
import org.apache.flink.api.common.typeinfo.Types;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.streaming.api.datastream.DataStreamSource;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.util.Collector;
import java.util.Collections;
import java.util.List;
/**
 * A demo that shows how to consume, convert, and then produce messages to ApsaraMQ for RocketMQ.
 */
public class RocketMQ5DataStreamDemo {

    public static final String ENDPOINT = "<yourEndpoint>";
    public static final String ACCESS_ID = "<accessID>";
    public static final String ACCESS_KEY = "<accessKey>";
    public static final String SOURCE_TOPIC = "<sourceTopicName>";
    public static final String CONSUMER_GROUP = "<consumerGroup>";
    public static final String SINK_TOPIC = "<sinkTopicName>";
    public static final String PRODUCER_GROUP = "<producerGroup>";

    public static void main(String[] args) throws Exception {
        // Set up the streaming execution environment
        Configuration conf = new Configuration();

        // The following two configurations are for local debugging only. Delete them before you package the job and upload it to Realtime Compute for Apache Flink.
        conf.setString("pipeline.classpaths", "file://" + "The absolute path of the uber JAR");
        conf.setString(
                "classloader.parent-first-patterns.additional",
                "com.alibaba.ververica.connectors.mq5.source.reader.deserializer.RocketMQRecordDeserializationSchema;com.alibaba.ververica.connectors.mq5.shaded.");
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(conf);

        final DataStreamSource<String> ds =
                env.fromSource(
                        RocketMQSource.<String>builder()
                                .setEndpoint(ENDPOINT)
                                .setAccessId(ACCESS_ID)
                                .setAccessKey(ACCESS_KEY)
                                .setTopic(SOURCE_TOPIC)
                                .setConsumerGroup(CONSUMER_GROUP)
                                .setDeserializationSchema(new MyDeserializer())
                                .setStartOffset(1)
                                .build(),
                        WatermarkStrategy.noWatermarks(),
                        "source");

        ds.map(new ToMessage())
                .addSink(
                        new OutputFormatSinkFunction<>(
                                new RocketMQOutputFormat.Builder()
                                        .setEndpoint(ENDPOINT)
                                        .setAccessId(ACCESS_ID)
                                        .setAccessKey(ACCESS_KEY)
                                        .setTopicName(SINK_TOPIC)
                                        .setProducerGroup(PRODUCER_GROUP)
                                        .build()));

        env.execute();
    }

    private static class MyDeserializer implements RocketMQRecordDeserializationSchema<String> {
        @Override
        public void deserialize(List<MessageExt> record, Collector<String> out) {
            for (MessageExt messageExt : record) {
                out.collect(new String(messageExt.getBody()));
            }
        }

        @Override
        public TypeInformation<String> getProducedType() {
            return Types.STRING;
        }
    }

    private static class ToMessage implements MapFunction<String, List<MessageExt>> {

        public ToMessage() {
        }

        @Override
        public List<MessageExt> map(String s) {
            final MessageExt message = new MessageExt();
            message.setBody(s.getBytes());
            message.setWaitStoreMsgOK(true);
            return Collections.singletonList(message);
        }
    }
}

RocketMQ 4.x

import com.alibaba.ververica.connector.mq.shaded.com.alibaba.rocketmq.common.message.MessageExt;
import com.alibaba.ververica.connectors.common.sink.OutputFormatSinkFunction;
import com.alibaba.ververica.connectors.metaq.sink.MetaQOutputFormat;
import com.alibaba.ververica.connectors.metaq.source.MetaQSource;
import com.alibaba.ververica.connectors.metaq.source.reader.deserializer.MetaQRecordDeserializationSchema;
import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.api.common.typeinfo.TypeInformation;
import org.apache.flink.api.connector.source.Boundedness;
import org.apache.flink.api.java.typeutils.ListTypeInfo;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.util.Collector;
import java.io.IOException;
import java.util.List;
import java.util.Properties;
import static com.alibaba.ververica.connector.mq.shaded.com.taobao.metaq.client.ExternConst.*;
/**
 * A demo that shows how to consume, convert, and then produce messages to ApsaraMQ for RocketMQ.
 */
public class RocketMQDataStreamDemo {

    public static final String ENDPOINT = "<yourEndpoint>";
    public static final String ACCESS_ID = "<accessID>";
    public static final String ACCESS_KEY = "<accessKey>";
    public static final String INSTANCE_ID = "<instanceID>";
    public static final String SOURCE_TOPIC = "<sourceTopicName>";
    public static final String CONSUMER_GROUP = "<consumerGroup>";
    public static final String SINK_TOPIC = "<sinkTopicName>";
    public static final String PRODUCER_GROUP = "<producerGroup>";

    public static void main(String[] args) throws Exception {
        // Set up the streaming execution environment
        Configuration conf = new Configuration();

        // The following two configurations are for local debugging only. Delete them before you package the job and upload it to Realtime Compute for Apache Flink.
        conf.setString("pipeline.classpaths", "file://" + "The absolute path of the uber JAR");
        conf.setString("classloader.parent-first-patterns.additional",
                "com.alibaba.ververica.connectors.metaq.source.reader.deserializer.MetaQRecordDeserializationSchema;com.alibaba.ververica.connector.mq.shaded.");
        final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(conf);

        // Create and add the ApsaraMQ for RocketMQ source.
        env.fromSource(createRocketMQSource(), WatermarkStrategy.noWatermarks(), "source")
                // Convert the message body to uppercase.
                .map(RocketMQDataStreamDemo2::convertMessages)
                // Create and add the ApsaraMQ for RocketMQ sink.
                .addSink(new OutputFormatSinkFunction<>(createRocketMQOutputFormat()))
                .name(RocketMQDataStreamDemo2.class.getSimpleName());
        // Compile and submit the job.
        env.execute("RocketMQ connector end-to-end DataStream demo");
    }

    private static MetaQSource<MessageExt> createRocketMQSource() {
        Properties mqProperties = createMQProperties();

        return new MetaQSource<>(SOURCE_TOPIC,
                CONSUMER_GROUP,
                null, // always null
                null, // tag of the messages to consume
                Long.MAX_VALUE, // stop timestamp in milliseconds
                -1, // start timestamp in milliseconds. Set to -1 to disable starting from an offset.
                0, // start offset
                300_000, // partition discovery interval
                mqProperties,
                Boundedness.CONTINUOUS_UNBOUNDED,
                new MyDeserializationSchema());
    }

    private static MetaQOutputFormat createRocketMQOutputFormat() {
        return new MetaQOutputFormat.Builder()
                .setTopicName(SINK_TOPIC)
                .setProducerGroup(PRODUCER_GROUP)
                .setMqProperties(createMQProperties())
                .build();
    }

    private static Properties createMQProperties() {
        Properties properties = new Properties();
        properties.put(PROPERTY_ONS_CHANNEL, "ALIYUN");
        properties.put(NAMESRV_ADDR, ENDPOINT);
        properties.put(PROPERTY_ACCESSKEY, ACCESS_ID);
        properties.put(PROPERTY_SECRETKEY, ACCESS_KEY);
        properties.put(PROPERTY_ROCKET_AUTH_ENABLED, true);
        properties.put(PROPERTY_INSTANCE_ID, INSTANCE_ID);
        return properties;
    }

    private static List<MessageExt> convertMessages(MessageExt messages) {
        return Collections.singletonList(messages);
    }

    public static class MyDeserializationSchema implements MetaQRecordDeserializationSchema<MessageExt> {
        @Override
        public void deserialize(List<MessageExt> list, Collector<MessageExt> collector) {
            for (MessageExt messageExt : list) {
                collector.collect(messageExt);
            }
        }

        @Override
        public TypeInformation<MessageExt> getProducedType() {
            return TypeInformation.of(MessageExt.class);
        }
    }
}
    }
}

XML

<!--ApsaraMQ for RocketMQ 5.x-->
<dependency>
    <groupId>com.alibaba.ververica</groupId>
    <artifactId>ververica-connector-mq5</artifactId>
    <version>${vvr-version}</version>
    <scope>provided</scope>
</dependency>

<!--ApsaraMQ for RocketMQ 4.x-->
<dependency>
    <groupId>com.alibaba.ververica</groupId>
    <artifactId>ververica-connector-mq</artifactId>
    <version>${vvr-version}</version>
</dependency>
Catatan

Untuk informasi selengkapnya tentang konfigurasi titik akhir untuk ApsaraMQ for RocketMQ, lihat Pengumuman tentang pengaturan titik akhir TCP internal.

FAQ

Bagaimana RocketMQ mendeteksi perubahan jumlah partisi selama penskalaan topik?