Topik ini memperkenalkan konektor ApsaraMQ for RocketMQ.
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 |
|
|
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 |
|
|
endPoint |
Titik akhir layanan. |
String |
Ya |
None |
ApsaraMQ for RocketMQ menyediakan dua jenis titik akhir:
Penting
Kami menyarankan Anda menggunakan titik akhir VPC. Koneksi jaringan publik mungkin tidak stabil karena perubahan dinamis pada kebijakan keamanan jaringan Alibaba Cloud.
|
|
topic |
Nama topik. |
String |
Ya |
None |
None |
|
accessId |
|
String |
|
None |
Penting
Untuk menghindari eksposur pasangan AccessKey Anda, kami menyarankan menggunakan variabel proyek untuk menentukan ID AccessKey dan Secret AccessKey.
|
|
accessKey |
|
String |
|
None |
|
|
tag |
Tag pesan yang akan berlangganan atau ditulis. |
String |
Tidak |
None |
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 |
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 |
|
lineDelimiter |
Pembatas baris yang digunakan untuk mengurai catatan. |
String |
Tidak |
\n |
None |
|
fieldDelimiter |
Pembatas bidang. |
String |
Tidak |
\u0001 |
Pembatas bervariasi berdasarkan mode terminal:
|
|
lengthCheck |
Kebijakan untuk memeriksa jumlah bidang dalam setiap catatan. |
String |
Tidak |
NONE |
Nilai yang valid:
|
|
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 Catatan
Didukung di VVR 8.0.5 dan versi yang lebih baru. |
|
deliveryTimestampMode |
Mode pengiriman untuk pesan tertunda. Parameter ini bekerja bersama parameter |
String |
Tidak |
None |
Nilai yang valid:
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:
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
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 |
String |
Tidak |
None |
Berlaku saat 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,femaleCatatanPesan 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>' );CatatanUntuk pesan RocketMQ dalam format biner, DDL harus mendefinisikan satu bidang dengan tipe data VARBINARY.
-
-
Buat tabel sink yang memetakan bidang metadata
keysdantagske 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
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
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 4.x: Konektor DataStream MQ.
-
ApsaraMQ for RocketMQ 5.x: Konektor DataStream MQ.
<!--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>
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?