DataHub は Apache Kafka プロトコルと完全に互換性があります。ネイティブ Kafka クライアントを使用して DataHub にデータを読み書きできます。
Kafka から DataHub へのマッピング
トピックタイプ
Kafka と DataHub では、トピックのスケーリングメカニズムが異なります。Kafka の動作との互換性を確保するため、DataHub トピックを作成する際は、スケーリングモードを ONLY_EXTEND に設定する必要があります。このモードでは、トピックに新しいシャードを追加することのみが可能です。このモードでは、シャードの分割やマージは許可されず、シャードの削除もまだサポートされていません。
トピック名
Kafka のトピック名は、DataHub のプロジェクトとトピックにマッピングされ、ピリオド (.) で区切られます。マッピングルールは次のとおりです。
最初の
.の前の部分が DataHub プロジェクト、その後の部分が DataHub トピックになります。たとえば、test_project.test_topicはtest_projectプロジェクトとtest_topicトピックにマッピングされます。名前に複数の
.が含まれる場合、最初の.のみが区切り文字として機能します。残りのすべての.および-は_に置き換えられます。
パーティション
DataHub の各アクティブシャードは、Kafka の 1 つのパーティションに対応します。たとえば、トピックに 5 つのアクティブシャードがある場合、5 つのパーティションを持つ Kafka トピックと同等になります。データを書き込む際、[0, 4] の範囲でパーティション ID を指定できます。パーティションを指定しない場合、Kafka クライアントが自動的に割り当てます。
Tuple トピック
Kafka から Tuple トピックにデータを書き込む場合、トピックスキーマには 1 つまたは 2 つの列が必要で、どちらの列も STRING 型である必要があります。それ以外の場合、書き込み操作は失敗します。
スキーマに列が 1 つの場合、値のみが書き込まれ、キーは破棄されます。
スキーマに列が 2 つある場合、1 列目と 2 列目がそれぞれキーと値に対応します。
また、Tuple トピックにバイナリデータを書き込むと文字化けが発生するため、書き込まないでください。バイナリデータを保存する場合は、Blob トピックを使用してください。
Blob トピック
Kafka から Blob トピックにデータを書き込む場合、Kafka メッセージの値が Blob フィールドに書き込まれます。メッセージキーが NULL でない場合、DataHub 属性として書き込まれます。属性名は __kafka_key__ で、その値は Kafka メッセージキーになります。
ヘッダー
Kafka ヘッダーは DataHub 属性にマッピングされます。ヘッダーの値が NULL の場合、無視され、属性として書き込まれません。Blob トピックの組み込み属性名との競合を避けるため、__kafka_key__ をヘッダーキーとして使用しないことを推奨します。
コンシューマーグループ
DataHub では、サブスクリプション ID がコンシューマーグループとして機能しますが、単一のトピックのみをサブスクライブできます。一方、Kafka コンシューマーグループは複数のトピックを同時にサブスクライブできます。Kafka のサブスクリプションモデルとの互換性を提供するため、DataHub はグループ機能を提供しています。プロジェクト内にグループを作成し、複数のトピックにバインドすることで、単一のグループですべてのトピックをサブスクライブできます。
グループは、サーバー上で複数の DataHub サブスクリプションを内部的に管理します。トピックをバインドすると、グループは自動的にサブスクリプションを作成し、トピックの詳細ページのサブスクリプションリストに表示されます。このサブスクリプションを手動で削除しないでください。削除すると、グループがトピックをサブスクライブできなくなり、既存のすべての消費オフセットが失われます。
1 つのグループは最大 50 個のトピックをサブスクライブできます。さらに多くのトピックをサブスクライブする場合は、チケットを提出してください。
Kafka パラメータ
C=コンシューマー、P=プロデューサー、S=Streams
パラメータ | C/P/S | 値 | 必須 | 説明 |
bootstrap.servers | * | 詳細については、「Kafka エンドポイント」をご参照ください。 | はい | |
security.protocol | * | SASL_SSL | はい | データセキュリティを確保するため、Kafka から DataHub への接続はデフォルトで SSL 暗号化を使用します。 |
sasl.mechanism | * | PLAIN | はい | AccessKey 認証情報の認証メカニズムです。PLAIN のみがサポートされています。 |
compression.type | P | LZ4 | いいえ | メッセージの圧縮タイプを指定します。現在、LZ4 のみがサポートされています。 |
group.id | C | project.topic:subId または project.group | はい |
|
partition.assignment.strategy | C | org.apache.kafka.clients.consumer.RangeAssignor | いいえ | Kafka のデフォルトのパーティション割り当て戦略は |
session.timeout.ms | C/S | [60000, 180000] | いいえ | Kafka のデフォルトは 10,000 ms です。ただし、DataHub では最小 60,000 ms が必要なため、この値は自動的に 60,000 ms に調整されます。 |
heartbeat.interval.ms | C/S | 推奨値:session.timeout.ms の 2/3 | いいえ | Kafka のデフォルトは 3,000 ms です。 |
application.id | S | project.topic:subId または project.group | はい |
|
この表は、Kafka クライアントを DataHub で使用する際に確認すべき主要なパラメータを示しています。retries や batch.size などの他のクライアント側パラメータは、ネイティブ Kafka と同様に動作します。サーバー側パラメータは DataHub の実際の動作を変更しません。たとえば、acks の値に関係なく、DataHub はデータが完全に書き込まれた後にのみ確認を返します。
Kafka エンドポイント
リージョン | リージョン ID | パブリックエンドポイント | ECS エンドポイント (クラシックネットワーク) | ECS エンドポイント (VPC) |
中国 (杭州) | cn-hangzhou | dh-cn-hangzhou.aliyuncs.com:9092 | dh-cn-hangzhou.aliyun-inc.com:9093 | dh-cn-hangzhou-int-vpc.aliyuncs.com:9094 |
中国 (上海) | cn-shanghai | dh-cn-shanghai.aliyuncs.com:9092 | dh-cn-shanghai.aliyun-inc.com:9093 | dh-cn-shanghai-int-vpc.aliyuncs.com:9094 |
中国 (北京) | cn-beijing | dh-cn-beijing.aliyuncs.com:9092 | dh-cn-beijing.aliyun-inc.com:9093 | dh-cn-beijing-int-vpc.aliyuncs.com:9094 |
中国 (張家口) | cn-zhangjiakou | dh-cn-zhangjiakou.aliyuncs.com:9092 | dh-cn-zhangjiakou.aliyun-inc.com:9093 | dh-cn-zhangjiakou-int-vpc.aliyuncs.com:9094 |
中国 (深圳) | cn-shenzhen | dh-cn-shenzhen.aliyuncs.com:9092 | dh-cn-shenzhen.aliyun-inc.com:9093 | dh-cn-shenzhen-int-vpc.aliyuncs.com:9094 |
シンガポール | ap-southeast-1 | dh-ap-southeast-1.aliyuncs.com:9092 | dh-ap-southeast-1.aliyun-inc.com:9093 | dh-ap-southeast-1-int-vpc.aliyuncs.com:9094 |
マレーシア (クアラルンプール) | ap-southeast-3 | dh-ap-southeast-3.aliyuncs.com:9092 | dh-ap-southeast-3.aliyun-inc.com:9093 | dh-ap-southeast-3-int-vpc.aliyuncs.com:9094 |
ドイツ (フランクフルト) | eu-central-1 | dh-eu-central-1.aliyuncs.com:9092 | dh-eu-central-1.aliyun-inc.com:9093 | dh-eu-central-1-int-vpc.aliyuncs.com:9094 |
中国東部 2 ファイナンス | cn-shanghai-finance-1 | dh-cn-shanghai-finance-1.aliyuncs.com:9092 | dh-cn-shanghai-finance-1.aliyun-inc.com:9093 | dh-cn-shanghai-finance-1-int-vpc.aliyuncs.com:9094 |
中国 (香港) | cn-hongkong | dh-cn-hongkong.aliyuncs.com:9092 | dh-cn-hongkong.aliyun-inc.com:9093 | dh-cn-hongkong-int-vpc.aliyuncs.com:9094 |
トピックの作成
コンソールでトピックを作成
トピックを作成する際、シャード拡張モードを有効にします。
SDK を使用してトピックを作成
Kafka API を使用してトピックを作成することはできません。DataHub SDK を使用し、
ExpandModeをONLY_EXTENDに設定する必要があります。必要な Maven 依存関係のバージョンは 2.19.0 以降です。AccessKey ID と AccessKey Secret の設定には環境変数を使用することを推奨します。プロジェクトコードにハードコーディングしないでください。Alibaba Cloud アカウントの AccessKey ペアはすべての API 操作に対する権限を持ちます。セキュリティ向上のため、API アクセスや日常業務には RAM ユーザーの AccessKey ペアを使用し、認証情報漏洩のリスクを軽減してください。
datahub.endpoint=<yourEndpoint> datahub.accessId=<yourAccessKeyId> datahub.accessKey=<yourAccessKeySecret><dependency> <groupId>com.aliyun.datahub</groupId> <artifactId>aliyun-sdk-datahub</artifactId> <version>2.19.0-public</version> </dependency>@Value("${datahub.endpoint}") String endpoint ; @Value("${datahub.accessId}") String accessId; @Value("${datahub.accessKey}") String accessKey; public class CreateTopic { public static void main(String[] args) { DatahubClient datahubClient = DatahubClientBuilder.newBuilder() .setDatahubConfig( new DatahubConfig(endpoint, new AliyunAccount(accessId, accessKey))) .build(); int shardCount = 1; int lifeCycle = 7; try { datahubClient.createTopic("test_project", "test_topic", shardCount, lifeCycle, RecordType.BLOB, "comment", ExpandMode.ONLY_EXTEND); } catch (DatahubClientException e) { e.printStackTrace(); } } }
グループの作成
コンソールでグループを作成
[グループの作成] をクリックし、右側のリストからサブスクライブするトピックを追加します。グループ作成後、バインドされたトピックを変更できます。グループは自動的にサブスクリプションを作成し、トピックのサブスクリプションリストページに表示されます。
SDK を使用してグループを作成
Maven 依存関係のバージョンは 2.21.6-public 以降である必要があります。
<dependency> <groupId>com.aliyun.datahub</groupId> <artifactId>aliyun-sdk-datahub</artifactId> <version>2.21.6-public</version> </dependency>@Value("${datahub.endpoint}") String endpoint ; @Value("${datahub.accessId}") String accessId; @Value("${datahub.accessKey}") String accessKey; public class CreateGroup { public static void main(String[] args) { DatahubClient datahubClient = DatahubClientBuilder.newBuilder() .setDatahubConfig( new DatahubConfig(endpoint, new AliyunAccount(accessId, accessKey))) .build(); List<String> topicList = new ArrayList<>(); topicList.add("test_project.topic1"); topicList.add("test_project.topic2"); topicList.add("test_project.topic3"); try { // Kafka グループを作成します。 datahubClient.createKafkaGroup("test_project", "test_topic", "test comment"); // サブスクリプション用にトピックをグループにバインドします。 datahubClient.updateTopicsForKafkaGroup("test_project", "test_topic", topicList, UpdateKafkaGroupMode.ADD); } catch (DatahubClientException e) { e.printStackTrace(); } } }
プロデューサーの例
kafka_client_producer_jaas.conf ファイル
任意のディレクトリに kafka_client_producer_jaas.conf という名前のファイルを作成し、次の内容を追加します。
KafkaClient {
org.apache.kafka.common.security.plain.PlainLoginModule required
username="yourAccessKeyId"
password="yourAccessKeySecret";
};Maven 依存関係
Kafka クライアントのバージョンは 0.10.0.0 以降である必要があります。バージョン 2.4.0 を推奨します。
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-clients</artifactId>
<version>2.4.0</version>
</dependency>サンプルコード
public class ProducerExample {
static {
System.setProperty("java.security.auth.login.config", "src/main/resources/kafka_client_producer_jaas.conf");
}
public static void main(String[] args) {
Properties properties = new Properties();
properties.put("bootstrap.servers", "dh-cn-hangzhou.aliyuncs.com:9092");
properties.put("security.protocol", "SASL_SSL");
properties.put("sasl.mechanism", "PLAIN");
properties.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
properties.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
properties.put("compression.type", "lz4");
String KafkaTopicName = "test_project.test_topic";
Producer<String, String> producer = new KafkaProducer<String, String>(properties);
try {
List<Header> headers = new ArrayList<>();
RecordHeader header1 = new RecordHeader("key1", "value1".getBytes());
RecordHeader header2 = new RecordHeader("key2", "value2".getBytes());
headers.add(header1);
headers.add(header2);
ProducerRecord<String, String> record = new ProducerRecord<>(KafkaTopicName, 0, "key", "Hello DataHub!", headers);
// 同期送信
producer.send(record).get();
} catch (InterruptedException e) {
e.printStackTrace();
} catch (ExecutionException e) {
e.printStackTrace();
} finally {
producer.close();
}
}
}結果
コードが正常に実行されると、データをサンプリングして結果を確認できます。
コンシューマーの例
kafka_client_producer_jaas.conf ファイルの生成方法および Maven 依存関係の追加方法については、プロデューサーの例をご参照ください。
新しいコンシューマーが参加すると、シャード割り当てに 10 ~ 20 秒かかります。割り当てが完了すると、コンシューマーはデータの消費を開始できます。
サンプルコード
Kafka グループの使用 (推奨)
package com.aliyun.datahub.kafka.demo;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import java.time.Duration;
import java.util.ArrayList;
import java.util.List;
import java.util.Properties;
public class ConsumerExample2 {
static {
System.setProperty("java.security.auth.login.config", "src/main/resources/kafka_client_producer_jaas.conf");
}
public static void main(String[] args) {
Properties properties = new Properties();
properties.put("bootstrap.servers", "dh-cn-hangzhou.aliyuncs.com:9092");
properties.put("security.protocol", "SASL_SSL");
properties.put("sasl.mechanism", "PLAIN");
// group.id を project.group 形式に設定します。
properties.put("group.id", "test_project.test_topic");
properties.put("auto.offset.reset", "earliest");
properties.put("session.timeout.ms", "60000");
properties.put("heartbeat.interval.ms", "40000");
properties.put("ssl.endpoint.identification.algorithm", "");
properties.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
properties.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
KafkaConsumer<String, String> kafkaConsumer = new KafkaConsumer<String, String>(properties);
List<String> topicList = new ArrayList<>();
topicList.add("test_project.test_topic1");
topicList.add("test_project.test_topic2");
topicList.add("test_project.test_topic3");
// Kafka グループを使用すると、複数のトピックをサブスクライブできます。
kafkaConsumer.subscribe(topicList);
while (true) {
ConsumerRecords<String, String> records = kafkaConsumer.poll(Duration.ofSeconds(5));
for (ConsumerRecord<String, String> record : records) {
System.out.println(record.toString());
}
}
}
}project.topic:subId の使用
package com.aliyun.datahub.kafka.demo;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import java.time.Duration;
import java.util.Collections;
import java.util.Properties;
public class ConsumerExample {
static {
System.setProperty("java.security.auth.login.config", "src/main/resources/kafka_client_producer_jaas.conf");
}
public static void main(String[] args) {
Properties properties = new Properties();
properties.put("bootstrap.servers", "dh-cn-hangzhou.aliyuncs.com:9092");
properties.put("security.protocol", "SASL_SSL");
properties.put("sasl.mechanism", "PLAIN");
// group.id を project.topic:subId 形式に設定します。
properties.put("group.id", "test_project.test_topic:1611039998153N71KM");
properties.put("auto.offset.reset", "earliest");
properties.put("session.timeout.ms", "60000");
properties.put("heartbeat.interval.ms", "40000");
properties.put("ssl.endpoint.identification.algorithm", "");
properties.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
properties.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
KafkaConsumer<String, String> kafkaConsumer = new KafkaConsumer<String, String>(properties);
// project.topic:subId 形式を使用する場合、単一のトピックのみをサブスクライブできます。
kafkaConsumer.subscribe(Collections.singletonList("test_project.test_topic"));
while (true) {
ConsumerRecords<String, String> records = kafkaConsumer.poll(Duration.ofSeconds(5));
for (ConsumerRecord<String, String> record : records) {
System.out.println(record.toString());
}
}
}
}結果
コードが正常に実行されると、ターミナルで消費されたデータを確認できます。 この例では、1 回のリクエストで返されるすべてのデータレコードは同じ LogAppendTime を持ちます。これは、そのバッチ内のすべてのレコードの最新のタイムスタンプです。
ConsumerRecord(topic = test_project.test_topic, partition = 0, leaderEpoch = 0, offset = 0, LogAppendTime = 1611040892661, serialized key size = 3, serialized value size = 14, headers = RecordHeaders(headers = [RecordHeader(key = key1, value = [118, 97, 108, 117, 101, 49]), RecordHeader(key = key2, value = [118, 97, 108, 117, 101, 50])], isReadOnly = false), key = key, value = Hello DataHub!)Streams の例
Maven 依存関係
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-clients</artifactId>
<version>2.4.0</version>
</dependency>
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-streams</artifactId>
<version>2.4.0</version>
</dependency>サンプルコード
この例では、test_project 内の入力トピックからデータを読み取り、キーと値の文字列を小文字に変換し、結果を出力トピックに書き込みます。
public class StreamExample {
static {
System.setProperty("java.security.auth.login.config", "src/main/resources/kafka_client_producer_jaas.conf");
}
public static void main(final String[] args) {
final String input = "test_project.input";
final String output = "test_project.output";
final Properties properties = new Properties();
properties.put("bootstrap.servers", "dh-cn-hangzhou.aliyuncs.com:9092");
properties.put("application.id", "test_project.input:1611293595417QH0WL");
properties.put("security.protocol", "SASL_SSL");
properties.put("sasl.mechanism", "PLAIN");
properties.put("session.timeout.ms", "60000");
properties.put("heartbeat.interval.ms", "40000");
properties.put("auto.offset.reset", "earliest");
final StreamsBuilder builder = new StreamsBuilder();
TestMapper testMapper = new TestMapper();
builder.stream(input, Consumed.with(Serdes.String(), Serdes.String()))
.map(testMapper)
.to(output, Produced.with(Serdes.String(), Serdes.String()));
final KafkaStreams streams = new KafkaStreams(builder.build(), properties);
final CountDownLatch latch = new CountDownLatch(1);
Runtime.getRuntime().addShutdownHook(new Thread("streams-shutdown-hook") {
@Override
public void run() {
streams.close();
latch.countDown();
}
});
try {
streams.start();
latch.await();
} catch (final Throwable e) {
System.exit(1);
}
System.exit(0);
}
static class TestMapper implements KeyValueMapper<String, String, KeyValue<String, String>> {
@Override
public KeyValue<String, String> apply(String s, String s2) {
return new KeyValue<>(StringUtils.lowerCase(s), StringUtils.lowerCase(s2));
}
}
}結果
Streams タスクを開始すると、シャード割り当てに約 1 分かかります。その後、コンソールで現在のタスク数を確認できます。タスク数は入力トピックのシャード数と一致します。この例では、入力トピックには 3 つのシャードがあります。
currently assigned active tasks: [0_0, 0_1, 0_2]
currently assigned standby tasks: []
revoked active tasks: []
revoked standby tasks: []シャードが割り当てられたら、(AAAA,BBBB),(CCCC,DDDD),(EEEE,FFFF) などのテストデータを入力トピックに書き込みます。次に、出力トピックからデータをサンプリングして、正しく書き込まれたことを確認します。
注意事項
トランザクションとべき等性はサポートされていません。
Kafka クライアントは DataHub でトピックを自動的に作成できません。データを書き込む前にトピックを作成する必要があります。
サブスクリプション ID (
project.topic:subid) をgroup.idとして使用する場合、コンシューマーは 1 つのトピックのみをサブスクライブできます。複数のトピックをサブスクライブするには、DataHub グループを使用してください。コンシューマーが読み取るデータのタイムスタンプは常に LogAppendTime であり、これはデータが DataHub に書き込まれた時刻を示します。1 回のフェッチリクエスト内のすべてのレコードは同じタイムスタンプを共有します。これは、そのバッチ内の最新のタイムスタンプです。つまり、読み取りタイムスタンプは実際の書き込み時刻よりも遅くなる可能性があります。
Streams アプリケーションは 1 つの入力トピックのみをサポートしますが、複数の出力トピックを持つことができます。
ステートレス Streams タスクのみがサポートされています。
サポートされている Kafka のバージョンは 0.10.0 ~ 2.4.0 です。
よくある質問
データ書き込み時に接続が切断される
Selector - [Producer clientId=producer-1] Connection with dh-cn-shenzhen.aliyuncs.com disconnected
java.io.EOFException
at org.apache.kafka.common.network.SslTransportLayer.read(SslTransportLayer.java:573)
...Kafka のメタデータリクエストとデータ書き込みリクエストは異なる接続を使用します。
クライアントは最初にメタデータを取得するための接続を確立します。次に、返されたブローカー情報を使用して、データ書き込み用の 2 番目の接続を確立します。以降のすべてのリクエストは、この 2 番目の接続を介して送信されます。
最初の接続はアイドル状態となり、タイムアウト後にサーバーによって自動的に閉じられます。これにより、ログに切断エラーが記録される場合があります。データが正常に書き込まれている場合は、このエラーを無視できます。
Kafka クライアントの起動に失敗する
Caused by: org.apache.kafka.common.errors.SslAuthenticationException: SSL handshake failed
Caused by: javax.net.ssl.SSLHandshakeException: No subject alternative names matching IP address 100.67.134.161 found設定に次のプロパティを追加してください。properties.put("ssl.endpoint.identification.algorithm", "");
消費中に DisconnectException が発生する
[INFO][Consumer clientId=client-id, groupId=consumer-project.topic:subid] Error sending fetch request (sessionId=INVALID, epoch=INITIAL) to node 1: {}.
org.apache.kafka.common.errors.DisconnectExceptionKafka クライアントはサーバーとの永続的な TCP 接続を維持する必要があります。この例外は通常、ネットワークの揺らぎが原因で発生します。クライアントには組み込みのリトライロジックがあるため、このエラーは通常、消費に影響を与えません。