このトピックでは、Flink から ClickHouse クラスターにデータをインポートする方法について説明します。
前提条件
-
Flink クラスターを作成済みであること。詳細については、「クラスターの作成」をご参照ください。
-
ClickHouse クラスターを作成済みであること。詳細については、「ClickHouse クラスターの作成」をご参照ください。
背景情報
Flink の詳細については、「Apache Flink」をご参照ください。
コード例
以下のセクションでは、ストリーム処理とバッチ処理のコード例を示します。
-
ストリーム処理
package com.company.packageName import java.util.concurrent.ThreadLocalRandom import scala.annotation.tailrec import org.apache.flink.api.common.typeinfo.Types import org.apache.flink.api.java.io.jdbc.JDBCAppendTableSink import org.apache.flink.streaming.api.scala._ import org.apache.flink.table.api.scala.{StreamTableEnvironment, table2RowDataStream} object StreamingJob { case class Test(id: Int, key1: String, value1: Boolean, key2: Long, value2: Double) private var dbName: String = "default" private var tableName: String = "" private var ckHost: String = "" private var ckPort: String = "8123" private var user: String = "default" private var password: String = "" def main(args: Array[String]) { parse(args.toList) checkArguments() // ストリーミング実行環境をセットアップします val env = StreamExecutionEnvironment.getExecutionEnvironment val tableEnv = StreamTableEnvironment.create(env) val insertIntoCkSql = s""" | INSERT INTO $tableName ( | id, key1, value1, key2, value2 | ) VALUES ( | ?, ?, ?, ?, ? | ) |""".stripMargin val jdbcUrl = s"jdbc:clickhouse://$ckHost:$ckPort/$dbName" println(s"jdbc url: $jdbcUrl") println(s"insert sql: $insertIntoCkSql") val sink = JDBCAppendTableSink .builder() .setDrivername("ru.yandex.clickhouse.ClickHouseDriver") .setDBUrl(jdbcUrl) .setUsername(user) .setPassword(password) .setQuery(insertIntoCkSql) .setBatchSize(1000) .setParameterTypes(Types.INT, Types.STRING, Types.BOOLEAN, Types.LONG, Types.DOUBLE) .build() val data: DataStream[Test] = env.fromCollection(1 to 1000).map(i => { val rand = ThreadLocalRandom.current() val randString = (0 until rand.nextInt(10, 20)) .map(_ => rand.nextLong()) .mkString("") Test(i, randString, rand.nextBoolean(), rand.nextLong(), rand.nextGaussian()) }) val table = table2RowDataStream(tableEnv.fromDataStream(data)) sink.emitDataStream(table.javaStream) // プログラムを実行します env.execute("Flink Streaming Scala API Skeleton") } private def printUsageAndExit(exitCode: Int = 0): Unit = { println("Usage: flink run com.company.packageName.StreamingJob /path/to/flink-clickhouse-demo-1.0.0.jar [options]") println(" --dbName ClickHouse データベースの名前を設定します。デフォルト: default。") println(" --tableName ClickHouse データベース内のテーブルの名前を設定します。") println(" --ckHost ClickHouse クラスターの IP アドレスを設定します。") println(" --ckPort ClickHouse クラスターのポートを設定します。デフォルト: 8123。") println(" --user ClickHouse クラスターにアクセスするためのユーザー名を設定します。") println(" --password ClickHouse ユーザーのパスワードを設定します。") System.exit(exitCode) } @tailrec private def parse(args: List[String]): Unit = args match { case ("--help" | "-h") :: _ => printUsageAndExit() case "--dbName" :: value :: tail => dbName = value parse(tail) case "--tableName" :: value :: tail => tableName = value parse(tail) case "--ckHost" :: value :: tail => ckHost = value parse(tail) case "--ckPort" :: value :: tail => ckPort = value parse(tail) case "--user" :: value :: tail => user = value parse(tail) case "--password" :: value :: tail => password = value parse(tail) case Nil => case _ => printUsageAndExit(1) } private def checkArguments(): Unit = { if ("".equals(tableName) || "".equals(ckHost)) { printUsageAndExit(2) } } } -
バッチ処理
package com.company.packageName import java.util.concurrent.ThreadLocalRandom import scala.annotation.tailrec import org.apache.flink.Utils import org.apache.flink.api.common.typeinfo.Types import org.apache.flink.api.java.io.jdbc.JDBCAppendTableSink import org.apache.flink.api.scala._ import org.apache.flink.table.api.scala.{BatchTableEnvironment, table2RowDataSet} object BatchJob { case class Test(id: Int, key1: String, value1: Boolean, key2: Long, value2: Double) private var dbName: String = "default" private var tableName: String = "" private var ckHost: String = "" private var ckPort: String = "8123" private var user: String = "default" private var password: String = "" def main(args: Array[String]) { parse(args.toList) checkArguments() // バッチ実行環境をセットアップします val env = ExecutionEnvironment.getExecutionEnvironment val tableEnv = BatchTableEnvironment.create(env) val insertIntoCkSql = s""" | INSERT INTO $tableName ( | id, key1, value1, key2, value2 | ) VALUES ( | ?, ?, ?, ?, ? | ) |""".stripMargin val jdbcUrl = s"jdbc:clickhouse://$ckHost:$ckPort/$dbName" println(s"jdbc url: $jdbcUrl") println(s"insert sql: $insertIntoCkSql") val sink = JDBCAppendTableSink .builder() .setDrivername("ru.yandex.clickhouse.ClickHouseDriver") .setDBUrl(jdbcUrl) .setUsername(user) .setPassword(password) .setQuery(insertIntoCkSql) .setBatchSize(1000) .setParameterTypes(Types.INT, Types.STRING, Types.BOOLEAN, Types.LONG, Types.DOUBLE) .build() val data = env.fromCollection(1 to 1000).map(i => { val rand = ThreadLocalRandom.current() val randString = (0 until rand.nextInt(10, 20)) .map(_ => rand.nextLong()) .mkString("") Test(i, randString, rand.nextBoolean(), rand.nextLong(), rand.nextGaussian()) }) val table = table2RowDataSet(tableEnv.fromDataSet(data)) sink.emitDataSet(Utils.convertScalaDatasetToJavaDataset(table)) // プログラムを実行します env.execute("Flink Batch Scala API Skeleton") } private def printUsageAndExit(exitCode: Int = 0): Unit = { println("Usage: flink run com.company.packageName.BatchJob /path/to/flink-clickhouse-demo-1.0.0.jar [options]") println(" --dbName ClickHouse データベースの名前を設定します。デフォルト: default。") println(" --tableName ClickHouse データベース内のテーブルの名前を設定します。") println(" --ckHost ClickHouse クラスターの IP アドレスを設定します。") println(" --ckPort ClickHouse クラスターのポートを設定します。デフォルト: 8123。") println(" --user ClickHouse クラスターにアクセスするためのユーザー名を設定します。") println(" --password ClickHouse ユーザーのパスワードを設定します。") System.exit(exitCode) } @tailrec private def parse(args: List[String]): Unit = args match { case ("--help" | "-h") :: _ => printUsageAndExit() case "--dbName" :: value :: tail => dbName = value parse(tail) case "--tableName" :: value :: tail => tableName = value parse(tail) case "--ckHost" :: value :: tail => ckHost = value parse(tail) case "--ckPort" :: value :: tail => ckPort = value parse(tail) case "--user" :: value :: tail => user = value parse(tail) case "--password" :: value :: tail => password = value parse(tail) case Nil => case _ => printUsageAndExit(1) } private def checkArguments(): Unit = { if ("".equals(tableName) || "".equals(ckHost)) { printUsageAndExit(2) } } }
手順
手順 1: ClickHouse テーブルを作成する
SSH モードで ClickHouse クラスターにログオンします。詳細については、クラスターへのログオンをご参照ください。
次のコマンドを実行して、ClickHouse クライアントを起動します。
clickhouse-client -h core-1-1 -m説明サンプルコマンドでは、core-1-1 はログオンするコアノードの名前を示します。複数のコアノードがある場合は、いずれかのノードにログオンできます。
ClickHouse データベースと必要な ClickHouse テーブルを作成します。
次のコマンドを実行して、clickhouse_database_name という名前のデータベースを作成します。
CREATE DATABASE clickhouse_database_name ON CLUSTER cluster_emr;Alibaba Cloud EMR は、cluster_emr という名前の ClickHouse クラスターを自動的に生成します。データベース名はカスタマイズできます。
次のコマンドを実行して、clickhouse_table_name_local という名前のテーブルを作成します。
CREATE TABLE clickhouse_database_name.clickhouse_table_name_local ON CLUSTER cluster_emr ( id UInt32, key1 String, value1 UInt8, key2 Int64, value2 Float64 ) ENGINE = ReplicatedMergeTree('/clickhouse/tables/{layer}-{shard}/clickhouse_database_name/clickhouse_table_name_local', '{replica}') ORDER BY id;説明テーブル名はカスタマイズできますが、テーブル名の末尾は _local にする必要があります。layer、shard、および replica パラメーターは、Alibaba Cloud EMR によって ClickHouse クラスター用に自動的に生成されるマクロであり、直接使用できます。
次のコマンドを実行して、clickhouse_table_name_all という名前のテーブルを作成します。このテーブルのフィールドは、clickhouse_table_name_local テーブルのフィールドと同じ方法で定義されています。
説明テーブル名はカスタマイズできますが、テーブル名の末尾は _all にする必要があります。
CREATE TABLE clickhouse_database_name.clickhouse_table_name_all ON CLUSTER cluster_emr ( id UInt32, key1 String, value1 UInt8, key2 Int64, value2 Float64 ) ENGINE = Distributed(cluster_emr, clickhouse_database_name, clickhouse_table_name_local, rand());
ステップ 2:コードのコンパイルとパッケージ化
-
flink-clickhouse-demo.tgz サンプルパッケージをローカルマシンにダウンロードして解凍します。
-
コマンドプロンプトで、pom.xml ファイルが含まれているディレクトリに移動し、次のコマンドを実行してファイルをパッケージ化します。
mvn clean packagepom.xml ファイルの artifactId に基づいて、target ディレクトリに flink-clickhouse-demo-1.0.0.jar JAR ファイルが生成されます。
ステップ 3:ジョブの送信
-
SSH を使用して Flink クラスターにログインします。詳細については、「クラスターへのログイン」をご参照ください。
-
パッケージ化された flink-clickhouse-demo-1.0.0.jar を Flink クラスターのルートディレクトリにアップロードします。
説明この例では、flink-clickhouse-demo-1.0.0.jar はルートディレクトリにアップロードされます。また、アップロードパスをカスタマイズすることもできます。
-
次のコマンドを実行してジョブを送信します。
コード例は次のとおりです:
-
ストリーム処理ジョブ
flink run -m yarn-cluster \ -c com.company.packageName.StreamingJob \ flink-clickhouse-demo-1.0.0.jar \ --dbName clickhouse_database_name \ --tableName clickhouse_table_name_all \ --ckHost ${clickhouse_host} \ --password ${password}; -
バッチ処理ジョブ
flink run -m yarn-cluster \ -c com.company.packageName.BatchJob \ flink-clickhouse-demo-1.0.0.jar \ --dbName clickhouse_database_name \ --tableName clickhouse_table_name_all \ --ckHost ${clickhouse_host} \ --password ${password};
パラメーター
説明
dbName
ClickHouse クラスターデータベースの名前。デフォルトは
defaultです。本トピックでは、clickhouse_database_name を例として使用します。tableName
ClickHouse クラスターデータベース内のテーブルの名前です。このトピックの例は clickhouse_table_name_all です。
ckHost
ClickHouse クラスターのマスターノードのプライベート IP アドレスまたはパブリック IP アドレス。IP アドレスの取得方法については、「マスターノードの IP アドレスの取得」をご参照ください。
password
ClickHouse ユーザーのパスワード。
パスワードは、ClickHouse サービスの Configure ページの users.default.password パラメーターから取得できます。
[configuration] タブで [server-users] カテゴリを選択すると、このパラメーターを確認できます。
-
マスターノードの IP アドレスの取得
[ノード] タブに移動します。
EMR コンソール にログオンします。左側のナビゲーションペインで、[ECS 上の EMR] をクリックします。
上部のナビゲーションバーで、クラスターが存在するリージョンとリソースグループをビジネス要件に基づいて選択します。
[ECS 上の EMR] ページで、管理するクラスターを見つけ、[アクション] 列の [ノード] をクリックします。
[ノード] タブで、マスターノードグループを見つけ、
アイコンをクリックし、[パブリック IP アドレス] 列の IP アドレスをコピーします。