ApsaraDB for SelectDB は Apache Doris と完全互換です。Flink Doris Connector を使用して、MySQL、Oracle、PostgreSQL、SQL Server、Kafka などのデータソースから既存データを SelectDB にインポートできます。Flink で変更データキャプチャ (CDC) タスクを開始すると、データソースからの増分データも SelectDB に同期されます。
概要
Flink Doris Connector は現在、SelectDB へのデータ書き込みのみをサポートしています。Flink Doris Connector を使用して SelectDB のバックエンドノードに直接接続し、効率的にデータを読み取る必要がある場合は、SelectDB テクニカルサポートチームに連絡してアクセス権をリクエストしてください。
SelectDB からデータを読み取るには、Flink JDBC Connector を使用することもできます。
Flink Doris Connector を使用すると、Flink が Apache Doris との間でリアルタイムのデータ処理および分析のためにデータの読み取りと書き込みを行えるようになります。SelectDB は Apache Doris と完全互換であるため、このコネクタは SelectDB へのストリーミングデータ投入の一般的な方法です。
各コンポーネントの機能は次のとおりです。
-
ソース
-
目的:ソースは外部システムから Flink データストリームにデータを読み取ります。これらのシステムには、メッセージキュー (Apache Kafka など)、データベース、ファイルシステムなどが含まれます。
-
例:Kafka をソースとしてリアルタイムメッセージを読み取る、またはファイルからデータを読み取る。
-
-
変換
-
目的:変換ステージでは、入力されたデータストリームを処理します。これらの操作には、フィルター、マッピング、集約、ウィンドウ処理などが含まれます。
-
例:入力ストリームをマッピングしてデータ構造を変換する、またはデータを集約して 1 分ごとのメトリックを計算する。
-
-
シンク
-
目的:シンクは、処理済みのデータを Flink データストリームから外部システム (データベース、ファイル、メッセージキューなど) に書き込みます。
-
例:処理結果を MySQL データベースに書き込む、またはデータを別の Kafka トピックに送信する。
-
次の図は、Flink Doris Connector を使用して SelectDB にデータをインポートする方法を示しています。
前提条件
-
ご利用のデータソース、Flink、および SelectDB 間のネットワーク接続を確保してください。
-
ApsaraDB for SelectDB インスタンスのパブリックエンドポイントを申請します。詳細については、「パブリックエンドポイントの申請またはリリース」をご参照ください。
ご利用の Flink 環境およびデータソースが ApsaraDB for SelectDB インスタンスと同じ VPC 内にある場合は、この手順をスキップできます。これは、それらが Alibaba Cloud プロダクトであるか、同じ VPC 内の Elastic Compute Service (ECS) インスタンス上にデプロイされている場合に一般的です。
-
ご利用の Flink 環境およびデータソースの IP アドレスを ApsaraDB for SelectDB インスタンスのホワイトリストに追加します。詳細については、「IP アドレスホワイトリストの設定」をご参照ください。
-
-
Flink Doris Connector がインストールされていることを確認してください。
次の表は、Flink および Flink Doris Connector のバージョン要件を示しています。
Flink バージョン
Flink Doris Connector バージョン
ダウンロードリンク
Realtime Compute for Apache Flink:1.17 以降
オープンソース Flink:1.15 以降
1.5.2 以降。最新バージョンのダウンロードを推奨します。
インストール手順については、「Flink Doris Connector のインストール」をご参照ください。
Flink Doris Connector の追加
ご利用の環境に基づいて Flink Doris Connector を追加します。
-
Realtime Compute for Apache Flink を使用して SelectDB にデータをインポートする場合は、
Flink Doris Connectorをカスタムコネクタとして管理できます。詳細については、「カスタムコネクタの管理」をご参照ください。 -
セルフマネージド Flink クラスターを使用する場合は、対応する
Flink Doris ConnectorJAR パッケージをダウンロードし、Flink インストールディレクトリのlibディレクトリに配置します。ダウンロードリンクについては、「JAR パッケージ」をご参照ください。 -
Flink Doris Connectorを Maven 依存関係として追加するには、プロジェクトの依存関係設定ファイルに次のコードを追加します。その他のバージョンについては、「Maven リポジトリ」をご参照ください。<!-- flink-doris-connector --> <dependency> <groupId>org.apache.doris</groupId> <artifactId>flink-doris-connector-1.16</artifactId> <version>1.5.2</version> </dependency>
例
例の環境
この例では、Flink SQL、Flink CDC、および DataStream API を使用して、ApsaraDB RDS for MySQL インスタンスの test データベースにある employees テーブルのデータを、SelectDB インスタンスの test データベースにある employees テーブルに移行します。これらの例のパラメーターを実際のシナリオに合わせて修正してください。例の環境は次のとおりです。
-
Flink 1.16 スタンドアロン環境
-
Java
-
ターゲットデータベース:test
-
ターゲットテーブル:employees
-
ソースデータベース:test
-
ソーステーブル:employees
環境の準備
Flink 環境
-
Java 環境を準備します。
Flink を実行するには Java 環境が必要です。Java 開発キット (JDK) をインストールし、
JAVA_HOME環境変数を設定する必要があります。サポートされている Java バージョンの一覧については、「Java Compatibility」をご参照ください。この例では「Java 8」を使用します。インストール手順については、「JDK のインストール」をご参照ください。
-
Flink インストールパッケージ flink-1.16.3-bin-scala_2.12.tgz をダウンロードします。このバージョンが古い場合は、「Apache Flink」から別のバージョンをダウンロードできます。
wget https://www.apache.si/flink/flink-1.16.3/flink-1.16.3-bin-scala_2.12.tgz -
インストールパッケージを解凍します。
tar -zxvf flink-1.16.3-bin-scala_2.12.tgz -
Flink インストールディレクトリの
libディレクトリに移動し、次の手順で必要なコネクタを追加します。-
Flink Doris Connector を追加します。
wget https://repo.maven.apache.org/maven2/org/apache/doris/flink-doris-connector-1.16/1.5.2/flink-doris-connector-1.16-1.5.2.jar -
Flink MySQL Connector を追加します。
wget https://repo1.maven.org/maven2/com/ververica/flink-sql-connector-mysql-cdc/2.4.2/flink-sql-connector-mysql-cdc-2.4.2.jar
-
-
Flink クラスターを起動します。
Flink インストールディレクトリの
binディレクトリで、次のコマンドを実行します。./start-cluster.sh
ターゲット SelectDB
-
ApsaraDB for SelectDB インスタンスを作成します。詳細については、「インスタンスの作成」をご参照ください。
-
インスタンスに接続します。詳細については、「インスタンスへの接続」をご参照ください。
-
testという名前のテストデータベースを作成します。CREATE DATABASE test; -
employeesという名前のテストテーブルを作成します。USE test; -- テーブルの作成 CREATE TABLE employees ( emp_no int NOT NULL, birth_date date, first_name varchar(20), last_name varchar(20), gender char(2), hire_date date ) UNIQUE KEY(`emp_no`) DISTRIBUTED BY HASH(`emp_no`) BUCKETS 1;
ソース MySQL
-
ApsaraDB RDS for MySQL インスタンス を作成します。
-
testという名前のテストデータベースを作成します。CREATE DATABASE test; -
employeesという名前のテストテーブルを作成します。USE test; CREATE TABLE employees ( emp_no INT NOT NULL PRIMARY KEY, birth_date DATE, first_name VARCHAR(20), last_name VARCHAR(20), gender CHAR(2), hire_date DATE ); -
データを挿入します。
INSERT INTO employees (emp_no, birth_date, first_name, last_name, gender, hire_date) VALUES (1001, '1985-05-15', 'John', 'Doe', 'M', '2010-06-20'), (1002, '1990-08-22', 'Jane', 'Smith', 'F', '2012-03-15'), (1003, '1987-11-02', 'Robert', 'Johnson', 'M', '2015-07-30'), (1004, '1992-01-18', 'Emily', 'Davis', 'F', '2018-01-05'), (1005, '1980-12-09', 'Michael', 'Brown', 'M', '2008-11-21');
Flink SQL を使用したインポート
-
Flink SQL クライアントを起動します。
Flink インストールディレクトリの
binディレクトリで、次のコマンドを実行します。./sql-client.sh -
Flink SQL クライアントで、Flink ジョブを送信します。
-
MySQL ソーステーブルを作成します。
次の文の
WITH句は、MySQL CDC Sourceの構成を指定します。パラメーターの詳細については、「MySQL | Apache Flink CDC」をご参照ください。CREATE TABLE employees_source ( emp_no INT, birth_date DATE, first_name STRING, last_name STRING, gender STRING, hire_date DATE, PRIMARY KEY (`emp_no`) NOT ENFORCED ) WITH ( 'connector' = 'mysql-cdc', 'hostname' = '127.0.0.1', 'port' = '3306', 'username' = 'root', 'password' = '****', 'database-name' = 'test', 'table-name' = 'employees' ); -
SelectDB 結果テーブルを作成します。
次の文の
WITH句は、SelectDB の構成を指定します。パラメーターの詳細については、「シンクパラメーター」をご参照ください。CREATE TABLE employees_sink ( emp_no INT , birth_date DATE, first_name STRING, last_name STRING, gender STRING, hire_date DATE ) WITH ( 'connector' = 'doris', 'fenodes' = 'selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080', 'table.identifier' = 'test.employees', 'username' = 'admin', 'password' = '****' ); -
MySQL ソーステーブルから SelectDB 結果テーブルにデータを同期します。
INSERT INTO employees_sink SELECT * FROM employees_source;
-
-
データインポートを検証します。
SelectDB に接続し、次の文を実行してインポートされたデータを表示します。
SELECT * FROM test.employees;
Flink CDC を使用したインポート
Realtime Compute for Apache Flink は JAR ベースのジョブをサポートしていません。代わりに CDC 3.0 で YAML ベースのジョブを使用してください。
Flink CDC を使用して SelectDB にデータをインポートします。
Flink CDC ジョブを実行するには、Flink インストールディレクトリの flink プログラムを使用します。構文は次のとおりです。
<FLINK_HOME>/bin/flink run \
-Dexecution.checkpointing.interval=10s \
-Dparallelism.default=1 \
-c org.apache.doris.flink.tools.cdc.CdcTools \
lib/flink-doris-connector-1.16-1.5.2.jar \
<mysql-sync-database|oracle-sync-database|postgres-sync-database|sqlserver-sync-database> \
--database <selectdb-database-name> \
[--job-name <flink-job-name>] \
[--table-prefix <selectdb-table-prefix>] \
[--table-suffix <selectdb-table-suffix>] \
[--including-tables <mysql-table-name|name-regular-expr>] \
[--excluding-tables <mysql-table-name|name-regular-expr>] \
--mysql-conf <mysql-cdc-source-conf> [--mysql-conf <mysql-cdc-source-conf> ...] \
--oracle-conf <oracle-cdc-source-conf> [--oracle-conf <oracle-cdc-source-conf> ...] \
--sink-conf <doris-sink-conf> [--table-conf <doris-sink-conf> ...] \
[--table-conf <selectdb-table-conf> [--table-conf <selectdb-table-conf> ...]]
パラメーター
|
パラメーター |
説明 |
|
execution.checkpointing.interval |
Flink チェックポイント間隔。この設定はデータ同期頻度に影響します。 |
|
parallelism.default |
Flink ジョブの並列度。並列度を上げるとデータ同期速度が向上します。 |
|
job-name |
Flink ジョブの名前。 |
|
database |
SelectDB のターゲットデータベース名。 |
|
table-prefix |
SelectDB のターゲットテーブル名のプレフィックス。例: |
|
table-suffix |
SelectDB のターゲットテーブル名のサフィックス。 |
|
including-tables |
同期するテーブル。複数のテーブルを区切るには縦棒 |
|
excluding-tables |
同期から除外するテーブル。 |
|
mysql-conf |
MySQL CDC Source の構成。詳細については、「MySQL CDC Connector」をご参照ください。 |
|
oracle-conf |
Oracle CDC Source の構成。詳細については、「Oracle CDC Connector」をご参照ください。 |
|
sink-conf |
Doris Sink の構成パラメーター。詳細については、「シンクパラメーター」をご参照ください。 |
|
table-conf |
SelectDB テーブルの構成パラメーター。これらは、SelectDB でテーブルを作成する際の |
-
データ同期を行うには、$FLINK_HOME/lib ディレクトリに必要な Flink CDC 依存関係(例:flink-sql-connector-mysql-cdc-${version}.jar または flink-sql-connector-oracle-cdc-${version}.jar)を追加する必要があります。
-
Flink 1.15 以降では、データベース全体の同期がサポートされています。Flink Doris Connector の異なるバージョンをダウンロードするには、「Flink Doris Connector」をご参照ください。
シンクパラメーター
|
パラメーター |
デフォルト |
必須 |
説明 |
|
fenodes |
なし |
はい |
ApsaraDB for SelectDB インスタンスのエンドポイントおよび HTTP ポート。 ApsaraDB for SelectDB コンソールの インスタンスの詳細 > ネットワーク情報 ページから、VPC アドレス(または インターネットアドレス)および HTTP プロトコルポート を取得できます。 例: |
|
table.identifier |
なし |
はい |
データベースおよびテーブル名。例: |
|
username |
なし |
はい |
ApsaraDB for SelectDB インスタンスのデータベースユーザー名。 |
|
password |
なし |
はい |
ApsaraDB for SelectDB インスタンスのデータベースユーザーのパスワード。 |
|
jdbc-url |
なし |
いいえ |
ApsaraDB for SelectDB インスタンスの JDBC 接続情報。 ApsaraDB for SelectDB コンソールの インスタンスの詳細 > ネットワーク情報 ページから、VPC アドレス(または インターネットアドレス)および MySQL プロトコルポート を取得できます。 例: |
|
auto-redirect |
true |
いいえ |
Stream Load リクエストのリダイレクトを有効にするかどうかを指定します。有効にすると、Stream Load はフロントエンド (FE) を介してデータを書き込み、バックエンド (BE) 情報は取得されません。 |
|
doris.request.retries |
3 |
いいえ |
SelectDB へのリクエスト送信のリトライ回数。 |
|
doris.request.connect.timeout |
30s |
いいえ |
SelectDB への接続タイムアウト。 |
|
doris.request.read.timeout |
30s |
いいえ |
SelectDB からのデータ読み取りタイムアウト。 |
|
sink.label-prefix |
"" |
はい |
Stream Load インポートのラベルプレフィックス。二相コミット (2PC) シナリオでは、Flink の 1 回限りのセマンティクス (EOS) を保証するために、このプレフィックスはグローバルに一意である必要があります。 |
|
sink.properties |
なし |
いいえ |
Stream Load のインポートパラメーター。プロパティは次のように構成します。
その他のパラメーターについては、「Stream Load」をご参照ください。 |
|
sink.buffer-size |
1048576 |
いいえ |
書き込みバッファのサイズ(バイト単位)。1 MB のデフォルト値を推奨します。 |
|
sink.buffer-count |
3 |
いいえ |
書き込みバッファの数。デフォルト値を推奨します。 |
|
sink.max-retries |
3 |
いいえ |
コミット失敗後の最大リトライ回数。デフォルトは 3 です。 |
|
sink.use-cache |
false |
いいえ |
例外発生時にメモリキャッシュを使用して復旧するかどうかを指定します。有効にすると、チェックポイント期間のデータがキャッシュに保持されます。 |
|
sink.enable-delete |
true |
いいえ |
削除イベントの同期を有効にするかどうかを指定します。このオプションは、Unique Key モデルを使用するテーブルでのみサポートされています。 |
|
sink.enable-2pc |
true |
いいえ |
二相コミット (2PC) を有効にするかどうかを指定します。1 回限りのセマンティクス (EOS) を保証するために、デフォルトで有効 ( |
|
sink.enable.batch-mode |
false |
いいえ |
SelectDB へのデータ書き込みにバッチモードを使用するかどうかを指定します。有効にすると、書き込み操作は Flink チェックポイントではなく、 バッチモードを有効にすると、1 回限りのセマンティクス (EOS) は保証されません。べき等性を実現するために Unique Key モデルを使用できます。 |
|
sink.flush.queue-size |
2 |
いいえ |
バッチモードでのバッファキューのサイズ。 |
|
sink.buffer-flush.max-rows |
50000 |
いいえ |
バッチモードでの 1 回のバッチ書き込みあたりの最大行数。 |
|
sink.buffer-flush.max-bytes |
10MB |
いいえ |
バッチモードでの 1 回のバッチ書き込みあたりの最大サイズ(バイト単位)。 |
|
sink.buffer-flush.interval |
10s |
いいえ |
バッチモードでの非同期バッファフラッシュ間隔。最小値は 1 秒です。 |
|
sink.ignore.update-before |
true |
いいえ |
|
同期の例
MySQL 同期
<FLINK_HOME>/bin/flink run \
-Dexecution.checkpointing.interval=10s \
-Dparallelism.default=1 \
-c org.apache.doris.flink.tools.cdc.CdcTools \
lib/flink-doris-connector-1.16-1.5.2.jar \
mysql-sync-database \
--database test \
--mysql-conf hostname=127.0.0.1 \
--mysql-conf port=3306 \
--mysql-conf username=root \
--mysql-conf password="password" \
--mysql-conf database-name=test \
--including-tables "employees" \
--sink-conf fenodes=selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080 \
--sink-conf username=admin \
--sink-conf password=****
Oracle 同期
<FLINK_HOME>/bin/flink run \
-Dexecution.checkpointing.interval=10s \
-Dparallelism.default=1 \
-c org.apache.doris.flink.tools.cdc.CdcTools \
lib/flink-doris-connector-1.16-1.5.2.jar \
oracle-sync-database \
--database test_db \
--oracle-conf hostname=127.0.0.1 \
--oracle-conf port=1521 \
--oracle-conf username=admin \
--oracle-conf password="password" \
--oracle-conf database-name=XE \
--oracle-conf schema-name=ADMIN \
--including-tables "tbl1|test.*" \
--sink-conf fenodes=selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080 \
--sink-conf username=admin \
--sink-conf password=****
PostgreSQL 同期
<FLINK_HOME>/bin/flink run \
-Dexecution.checkpointing.interval=10s \
-Dparallelism.default=1 \
-c org.apache.doris.flink.tools.cdc.CdcTools \
lib/flink-doris-connector-1.16-1.5.2.jar \
postgres-sync-database \
--database db1\
--postgres-conf hostname=127.0.0.1 \
--postgres-conf port=5432 \
--postgres-conf username=postgres \
--postgres-conf password="123456" \
--postgres-conf database-name=postgres \
--postgres-conf schema-name=public \
--postgres-conf slot.name=test \
--postgres-conf decoding.plugin.name=pgoutput \
--including-tables "tbl1|test.*" \
--sink-conf fenodes=selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080 \
--sink-conf username=admin \
--sink-conf password=****
SQL Server 同期
<FLINK_HOME>/bin/flink run \
-Dexecution.checkpointing.interval=10s \
-Dparallelism.default=1 \
-c org.apache.doris.flink.tools.cdc.CdcTools \
lib/flink-doris-connector-1.16-1.5.2.jar \
sqlserver-sync-database \
--database db1\
--sqlserver-conf hostname=127.0.0.1 \
--sqlserver-conf port=1433 \
--sqlserver-conf username=sa \
--sqlserver-conf password="123456" \
--sqlserver-conf database-name=CDC_DB \
--sqlserver-conf schema-name=dbo \
--including-tables "tbl1|test.*" \
--sink-conf fenodes=selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080 \
--sink-conf username=admin \
--sink-conf password=****
DataStream API を使用したインポート
-
Maven プロジェクトに次の依存関係を追加します。
Maven 依存関係
-
コア Java コード。
次のコードは、MySQL ソーステーブルおよび ApsaraDB for SelectDB 結果テーブルを構成します。パラメーターは、「Flink SQL を使用したデータインポート」セクションで使用したものに対応しています。詳細については、「MySQL | Apache Flink CDC」および「シンクパラメーター」をご参照ください。
package org.example; import com.ververica.cdc.connectors.mysql.source.MySqlSource; import com.ververica.cdc.connectors.mysql.table.StartupOptions; import com.ververica.cdc.connectors.shaded.org.apache.kafka.connect.json.JsonConverterConfig; import com.ververica.cdc.debezium.JsonDebeziumDeserializationSchema; import org.apache.doris.flink.cfg.DorisExecutionOptions; import org.apache.doris.flink.cfg.DorisOptions; import org.apache.doris.flink.sink.DorisSink; import org.apache.doris.flink.sink.writer.serializer.JsonDebeziumSchemaSerializer; import org.apache.doris.flink.tools.cdc.mysql.DateToStringConverter; import org.apache.flink.api.common.eventtime.WatermarkStrategy; import org.apache.flink.streaming.api.datastream.DataStreamSource; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import java.util.HashMap; import java.util.Map; import java.util.Properties; public class Main { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); env.setParallelism(1); env.enableCheckpointing(10000); Map<String, Object> customConverterConfigs = new HashMap<>(); customConverterConfigs.put(JsonConverterConfig.DECIMAL_FORMAT_CONFIG, "numeric"); JsonDebeziumDeserializationSchema schema = new JsonDebeziumDeserializationSchema(false, customConverterConfigs); // MySQL ソーステーブルの構成 MySqlSource<String> mySqlSource = MySqlSource.<String>builder() .hostname("rm-xxx.mysql.rds.aliyuncs***") .port(3306) .startupOptions(StartupOptions.initial()) .databaseList("db_test") .tableList("db_test.employees") .username("root") .password("test_123") .debeziumProperties(DateToStringConverter.DEFAULT_PROPS) .deserializer(schema) .serverTimeZone("Asia/Shanghai") .build(); // ApsaraDB for SelectDB 結果テーブルの構成 DorisSink.Builder<String> sinkBuilder = DorisSink.builder(); DorisOptions.Builder dorisBuilder = DorisOptions.builder(); dorisBuilder.setFenodes("selectdb-cn-xxx-public.selectdbfe.rds.aliyunc****:8080") .setTableIdentifier("db_test.employees") .setUsername("admin") .setPassword("test_123"); DorisOptions dorisOptions = dorisBuilder.build(); // sink.properties を使用した Stream Load パラメーターの構成 Properties properties = new Properties(); properties.setProperty("format", "json"); properties.setProperty("read_json_by_line", "true"); DorisExecutionOptions.Builder executionBuilder = DorisExecutionOptions.builder(); executionBuilder.setStreamLoadProp(properties); sinkBuilder.setDorisExecutionOptions(executionBuilder.build()) .setSerializer(JsonDebeziumSchemaSerializer.builder().setDorisOptions(dorisOptions).build()) // データストリームをシリアル化。 .setDorisOptions(dorisOptions); DataStreamSource<String> dataStreamSource = env.fromSource(mySqlSource, WatermarkStrategy.noWatermarks(), "MySQL Source"); dataStreamSource.sinkTo(sinkBuilder.build()); env.execute("MySQL to SelectDB"); } }
高度な使用方法
Flink SQL を使用した部分カラムの更新
-- チェックポイントを有効化
SET 'execution.checkpointing.interval' = '10s';
CREATE TABLE cdc_mysql_source (
id INT
,name STRING
,bank STRING
,age INT
,PRIMARY KEY (id) NOT ENFORCED
) WITH (
'connector' = 'mysql-cdc',
'hostname' = '127.0.0.1',
'port' = '3306',
'username' = 'root',
'password' = 'password',
'database-name' = 'database',
'table-name' = 'table'
);
CREATE TABLE selectdb_sink (
id INT,
name STRING,
bank STRING,
age INT
)
WITH (
'connector' = 'doris',
'fenodes' = 'selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080',
'table.identifier' = 'database.table',
'username' = 'admin',
'password' = '****',
'sink.properties.format' = 'json',
'sink.properties.read_json_by_line' = 'true',
'sink.properties.columns' = 'id,name,bank,age',
'sink.properties.partial_columns' = 'true' -- 部分カラム更新を有効化。
);
INSERT INTO selectdb_sink SELECT id,name,bank,age FROM cdc_mysql_source;
Flink SQL を使用して 列単位のデータ を削除する
CDC シナリオでは、Doris シンクは RowKind からイベントタイプを識別し、隠しカラム __DORIS_DELETE_SIGN__ に値を割り当てて削除を実行します。データソースが Kafka メッセージの場合、シンクは RowKind を使用して操作タイプを判断できません。代わりに、{"op_type":"delete",data:{...}} のようなメッセージ内の特定のフィールドに依存する必要があります。Kafka データの特定フィールドに基づいて Alibaba Cloud SelectDB のデータを削除する方法を次の Flink SQL の例に示します。
-- 例のメッセージ: {"op_type":"delete",data:{"id":1,"name":"zhangsan"}}
CREATE TABLE KAFKA_SOURCE(
data STRING,
op_type STRING
) WITH (
'connector' = 'kafka',
...
);
CREATE TABLE SELECTDB_SINK(
id INT,
name STRING,
__DORIS_DELETE_SIGN__ INT
) WITH (
'connector' = 'doris',
'fenodes' = 'selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080',
'table.identifier' = 'db.table',
'username' = 'admin',
'password' = '****',
'sink.enable-delete' = 'false', -- false の値は、イベントタイプが RowKind から推論されないことを示します。
'sink.properties.columns' = 'id, name, __DORIS_DELETE_SIGN__' -- Stream Load インポートのカラムを明示的に指定します。
);
INSERT INTO SELECTDB_SINK
SELECT json_value(data,'$.id') as id,
json_value(data,'$.name') as name,
if(op_type='delete',1,0) as __DORIS_DELETE_SIGN__
FROM KAFKA_SOURCE;
よくある質問
-
Q:
BITMAPデータを書き込むにはどうすればよいですか?A:次の例をご参照ください。
CREATE TABLE bitmap_sink ( dt INT, page STRING, user_id INT ) WITH ( 'connector' = 'doris', 'fenodes' = 'selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080', 'table.identifier' = 'test.bitmap_test', 'username' = 'admin', 'password' = '****', 'sink.label-prefix' = 'selectdb_label', 'sink.properties.columns' = 'dt,page,user_id,user_id=to_bitmap(user_id)' ); -
Q:
errCode = 2, detailMessage = Label[label_0_1]has already been used, relate to txn[19650]エラーを解決するにはどうすればよいですか?A:1 回限りのセマンティクスのシナリオでは、Flink ジョブを最新のチェックポイントまたはセーブポイントから再起動する必要があります。古い状態からジョブを再起動すると、このエラーが発生します。1 回限りのセマンティクスが不要な場合は、
sink.enable-2pc=falseを設定して二相コミット (2PC) を無効にするか、別の sink.label-prefix を使用してください。 -
Q:
errCode = 2, detailMessage = transaction[19650]not foundエラーを解決するにはどうすればよいですか?A:このエラーはコミットフェーズ中に発生します。チェックポイントに記録されたトランザクション ID が ApsaraDB for SelectDB で期限切れになっていることを示しています。コネクタがこの期限切れのトランザクションをコミットしようとすると、サーバーはトランザクションが見つからないと報告します。この場合、チェックポイントからジョブを再起動できません。この問題を防ぐには、ApsaraDB for SelectDB の
streaming_label_keep_max_secondパラメーターを増やしてください。デフォルト値は 12 時間です。 -
Q:
errCode = 2, detailMessage = current running txns on db 10006 is 100, larger than limit 100エラーを解決するにはどうすればよいですか?A:このエラーは、単一データベースの同時インポートトランザクション数がシステム制限の 100 を超えていることを示しています。これを解決するには、ApsaraDB for SelectDB の
max_running_txn_num_per_dbパラメーターを増やしてください。詳細については、「max_running_txn_num_per_db」をご参照ください。このエラーは、ラベルを頻繁に変更してジョブを再起動する場合にも発生する可能性があります。二相コミット (2PC) シナリオ(Duplicate Key および Aggregate Key モデルに適用)では、各ジョブに一意のラベルが必要です。チェックポイントからジョブを再起動すると、Flink は事前コミット済みだがまだコミットされていないトランザクションのみを中止します。再起動前にラベルを頻繁に変更すると、多くの事前コミット済みトランザクションが中止されず、トランザクションクォータを消費し続けます。Unique Key モデルの場合は、2PC を無効にして、べき等な書き込みを行うようにシンクオペレーターを設計できます。
-
Q:Unique Key モデルを使用するテーブルに書き込む際に、バッチ内のデータ順序を保証するにはどうすればよいですか?
A:データ順序を保証するためにシーケンスカラム構成を追加します。詳細については、「SEQUENCE」をご参照ください。
-
Q:Flink ジョブにエラーが報告されていないのに、データが同期されないのはなぜですか?
A:この動作はコネクタのバージョンによって異なります。1.1.0 より前のバージョンでは、書き込みはバッチ処理され、データ駆動型であるため、アップストリームソースがデータを生成していることを確認する必要があります。1.1.0 以降のバージョンでは、書き込みはチェックポイントによってトリガーされるため、データを書き込むにはチェックポイントを有効にする必要があります。
-
Q:
tablet writer write failed, tablet_id=190958, txn_id=3505530, err=-235エラーを解決するにはどうすればよいですか?A:このエラーは通常、1.1.0 より前のコネクタバージョンで発生します。書き込み頻度が高すぎることで、タブレットにバージョンが多数作成されることが原因です。この問題を解決するには、
sink.buffer-flush.max-bytesおよびsink.buffer-flush.intervalパラメーターを増やして、Stream Load の頻度を下げてください。 -
Q:Flink インポート中にダーティデータをスキップするにはどうすればよいですか?
A:ソースデータに送信先テーブルのスキーマと一致しないレコード(例:データ型または長さが不正)が含まれている場合、Stream Load ジョブが失敗し、Flink が継続的にリトライします。このダーティデータをスキップするには、Stream Load の厳格モードを無効にして
strict_mode=false,max_filter_ratio=1を設定するか、シンクオペレーターに到達する前に無効なデータをフィルターで除外する変換ステップを追加します。 -
Q:ソーステーブルを ApsaraDB for SelectDB テーブルにどのようにマッピングすればよいですか?
A:Flink Doris Connector を使用してデータをインポートする際は、次のマッピングが正しいことを確認してください。(1) ソーステーブルのカラムおよび型が Flink SQL と一致していること。(2) Flink SQL のカラムおよび型が ApsaraDB for SelectDB テーブルと一致していること。
-
Q:
TApplicationException: get_next failed: out of sequence response: expected 4 but got 3エラーを解決するにはどうすればよいですか?A:このエラーは、基盤となる Thrift フレームワーク内の同時実行バグを示しています。これを解決するには、Flink Doris Connector の最新バージョンにアップグレードし、互換性のある Flink バージョンを使用してください。
-
Q:
DorisRuntimeException: Fail to abort transaction 26153 with urlhttp://192.168.XX.XXエラーを解決するにはどうすればよいですか?A:この問題を診断するには、TaskManager ログで
abort transaction responseというフレーズを検索します。ログエントリの HTTP ステータスコードにより、問題がクライアント側またはサーバー側のどちらに起因するかを判断できます。