本トピックでは、DataHub の新しいバッチ処理機能について、その仕組みと実装方法を説明します。バッチモードに切り替えることで、サーバーサイドのリソース消費が減り、使用料金が大幅に削減されるため、パフォーマンスが向上し、コストが削減されます。
アップグレードの詳細
zstd 圧縮のサポート
DataHub は、Zstandard (zstd) 圧縮アルゴリズムをサポートするようになりました。DataHub がサポートしている lz4 および deflate アルゴリズムと比較して、zstd は優れたパフォーマンスを発揮します。
Zstandard (zstd) は、Facebook によって開発され、2016 年にオープンソース化された高性能な圧縮アルゴリズムです。zstd アルゴリズムは圧縮速度と圧縮率の両方に優れており、DataHub のシナリオに非常に適しています。
シリアル化の変換
DataHub は、データを転送用に構成する方式であるバッチシリアル化を導入しました。「バッチ」という用語は、特定のシリアル化方式を指すものではありません。代わりに、シリアル化されたデータの二次的なカプセル化を指します。たとえば、一度に 100 レコードを送信するには、まず 100 レコードをバッファーにシリアル化します。次に、このバッファーに圧縮アルゴリズムを適用して、圧縮済みバッファーを作成します。最後に、ヘッダーを圧縮済みバッファーに追加して、サイズ、データレコード数、圧縮アルゴリズム、CRC チェックサムなどのメタデータを格納します。このプロセスにより、完全なバッチバッファーが作成されます。
これにより、以下の問題が解決されます:
ビジネスレイヤーでのダーティデータを効果的に防止します。
サーバーサイドの CPU オーバーヘッドを削減し、データ処理性能を向上させます。
同時読み書きのレイテンシーを削減します。
クライアントはバッチバッファーを送信する前に完全なデータ有効性チェックを実行するため、サーバーは CRC チェックサムを検証するだけでバッファーが有効であることを確認できます。その後、バッファーは直接ディスクに書き込むことができます。このプロセスにより、サーバーサイドのシリアル化、デシリアライズ、圧縮、展開、および検証操作が不要になります。その結果、サーバーサイドのパフォーマンスは 80% 以上向上します。複数のデータレコードがまとめて圧縮されるため、圧縮率も向上し、ストレージコストが削減されます。
コストの比較
バッチ処理の効果を検証するために、以下のテストが実施されました。テストシナリオは次のとおりです:
テストデータは、約 200 列の広告配信情報で構成されています。データ値の約 20% から 30% が null です。
バッチあたり 1,000 データレコード。
バッチ内のシリアル化には Avro が使用されます。
以前のバージョンのデフォルトの圧縮アルゴリズムである lz4 の代わりに、zstd 圧縮アルゴリズムが使用されます。
テストでは、以下の結果が得られました:
データソースサイズ (バイト) | lz4 圧縮サイズ (バイト) | zstd 圧縮サイズ (バイト) | |
protobuf シリアル化 | 11,506,677 | 3,050,640 | 1,158,868 |
バッチシリアル化 | 11,154,596 | 2,931,729 | 1,112,693 |
DataHub の課金は、主にストレージとトラフィックの 2 つの側面に基づいています。他の課金項目は、主にリソースの乱用を防ぐためのペナルティとして意図されています。したがって、テスト結果は ストレージ と トラフィック の観点から分析します。
ストレージコスト:protobuf シリアル化を使用する場合、DataHub はストレージ用にデータを圧縮しません。データは HTTP 送信中にのみ圧縮されます。この方式をバッチ処理と zstd 圧縮に置き換えると、ストレージサイズは 11,506 KB から 1,112 KB に減少します。これにより、ストレージコストが約 90% 削減されます。
トラフィックコスト:protobuf と lz4 を使用する場合、データサイズは 3,050 KB です。バッチ処理と zstd を使用する場合、データサイズは 1,112 KB です。これにより、トラフィックコストが約 60% 削減されます。
前述の結果はサンプルデータに基づいています。実際の結果はデータによって異なる場合があります。必要に応じて独自のテストを実行することを推奨します。
バッチ処理の使用
注意事項
バッチ書き込みの主な利点は、各バッチのレコード数を最大化することです。クライアントがレコードをバッチ処理できない場合、またはバッチあたりのレコード数が少ない場合、パフォーマンスの向上は大きくない可能性があります。
円滑な移行を確実にするため、バッチモードは元の読み書きメソッドと互換性があります。バッチモードで書き込まれたデータは元のモードで読み取ることができ、その逆も同様です。ただし、最適なパフォーマンスを得るには、同じモードでデータを書き込み、消費する必要があります。 [書き込みと消費のモードが一致しないと、パフォーマンスが低下します]。
DataHub の最新バージョンでは、バッチプロトコルを使用するためにトピックのマルチバージョン スキーマを有効にする必要はなくなりました。ただし、クライアントをサポートされているバージョンに更新する必要があります。次の表に、サポートされているクライアントと必要なバージョンを示します。
言語 | サポート状況 | バージョン要件 |
Java SDK | サポート済み | 1.5.1 以降 |
Go SDK | サポート済み | 1.1.0 以降 |
Python | 未サポート | - |
C++ | 未サポート | - |
Java の例
バージョン 1.5.1 以降では、デフォルトのシリアル化プロトコルは Batch です。特別な設定は不要です。同期書き込みまたは非同期書き込みを使用して、通常どおりデータを書き込むことができます。
Maven 依存関係
<dependency>
<groupId>com.aliyun.datahub</groupId>
<artifactId>datahub-client-library</artifactId>
<version>1.5.1</version>
</dependency>コード例
public static void main(String[] args) throws InterruptedException {
// 環境変数から AccessKey 情報を取得します。
EnvironmentVariableCredentialProvider provider = EnvironmentVariableCredentialProvider.create();
String endpoint ="https://dh-cn-hangzhou.aliyuncs.com";
String projectName = "test_project";
String topicName = "test_topic";
// プロデューサーを初期化します。デフォルト設定が使用されます。
ProducerConfig config = new ProducerConfig(endpoint, provider);
DatahubProducer producer = new DatahubProducer(projectName, topicName, config);
RecordSchema schema = producer.getTopicSchema();
// マルチバージョン スキーマが有効な場合は、指定したバージョンのスキーマも取得できます。
// RecordSchema schema = producer.getTopicSchema(3);
// 非同期書き込みの場合、必要に応じてコールバック関数を登録します。
WriteCallback callback = new WriteCallback() {
@Override
public void onSuccess(String shardId, List<RecordEntry> records, long elapsedTimeMs, long sendTimeMs) {
System.out.println("write success");
}
@Override
public void onFailure(String shardId, List<RecordEntry> records, long elapsedTimeMs, DatahubClientException e) {
System.out.println("write failed");
}
};
for (int i = 0; i < 10000; ++i) {
try {
// スキーマに従ってデータを生成します。
TupleRecordData data = new TupleRecordData(schema);
data.setField("field1", "hello");
data.setField("field2", 1234);
RecordEntry recordEntry = new RecordEntry();
recordEntry.setRecordData(data);
producer.sendAsync(recordEntry, callback);
// データが正常に送信されたかを確認する必要がない場合は、コールバックを登録せずに直接データを送信します。
// producer.sendAsync(recordEntry, null);
} catch (DatahubClientException e) {
// TODO: 例外を処理します。通常は再試行不可能なエラー、または最大再試行回数に達した後のエラーです。
Thread.sleep(1000);
}
}
// プログラムが終了する前に、すべてのデータが送信されるようにします。
producer.flush(true);
producer.close();
}Go の例
go.mod 依存関係
require (
github.com/aliyun/aliyun-datahub-sdk-go v1.1.0
)コード例
バージョン 1.1.0 以降では、デフォルトのシリアル化プロトコルは Batch です。特別な設定は不要です。同期書き込みまたは非同期書き込みを使用して、通常どおりデータを書き込むことができます。
func handleSuccessRun(producer datahub.AsyncProducer) {
for suc := range producer.Successes() {
// リクエスト成功の処理
fmt.Printf("shard:%s, rid:%s, records:%d, latency:%v\n",
suc.ShardId, suc.RequestId, len(suc.Records), suc.Latency)
}
}
func handleFailedRun(producer datahub.AsyncProducer) {
// リクエスト失敗の処理
for err := range producer.Errors() {
fmt.Printf("shard:%s, records:%d, latency:%v, error:%v\n",
err.ShardId, len(err.Records), err.Latency, err.Err)
}
}
func main() {
cfg := datahub.NewProducerConfig()
// 環境変数から AccessKey 情報を取得します。
credential, err := credentials.NewCredential(nil)
if err != nil {
fmt.Println(err)
// TODO: エラー処理
}
cfg.Account = datahub.NewCredentialAccount(credential)
cfg.Endpoint = "https://dh-cn-hangzhou.aliyuncs.com"
cfg.Project = "test_project"
cfg.Topic = "test_topic"
producer := datahub.NewAsyncProducer(cfg)
err = producer.Init()
if err != nil {
// TODO: エラー処理
fmt.Println(err)
}
schema, err := producer.GetSchema()
if err != nil {
// TODO: エラー処理
fmt.Println(err)
}
// 成功チャネルの処理
go handleSuccessRun(producer)
// エラーチャネルの処理
go handleFailedRun(producer)
// ループでデータを生成します。
for i := 0; i < 1000; i++ {
record := datahub.NewTupleRecord(schema)
// 値を設定するたびに、操作が成功したかどうかを確認します。
err = record.SetValueByName("f1", "val1")
if err != nil {
fmt.Println(err)
return
}
err = record.SetValueByName("f2", 1234)
if err != nil {
fmt.Println(err)
return
}
producer.Input() <- record
}
err = producer.Close()
if err != nil {
fmt.Println(err)
}
}
サポート
ご質問や問題がある場合は、チケットを起票してサポートにお問い合わせいただくか、ユーザーグループにご参加ください。グループ番号は 33517130 です。