Amazon Glue は、Amazon Web Services (AWS) が提供する完全マネージド型のサーバーレス データインテグレーションサービスです。分析、機械学習、アプリケーション開発に向けて、データの検出、準備、変換を簡素化します。
インストール
Glue コードを ClickHouse と統合するには、次のいずれかの方法で Glue から公式の Spark コネクタを使用できます。
- AWS Marketplace から ClickHouse Glue コネクタをインストールする (推奨) 。
- Spark Connector の JAR を Glue ジョブに手動で追加する。
コネクタをサブスクライブする
アカウントでコネクタを利用するには、AWS Marketplace で ClickHouse AWS Glue Connector をサブスクライブします。
必要な権限を付与する
最小権限のガイドに記載されているとおり、Glue ジョブの IAM ロールに必要な権限があることを確認してください。
コネクタを有効化して接続を作成する
サブスクライブ後、ジョブの要件に合った Glue バージョンを選択します。Additional details セクションの Usage instructions で、Open Glue Studio - Add ClickHouse connector へのリンクをクリックします。すると、主要なフィールドが事前入力された Glue の接続作成ページが開きます。接続に名前を付けて、作成をクリックしてください (この段階では ClickHouse の接続情報を指定する必要はありません) 。

Glue ジョブで使用する
Glue ジョブで Job details タブを選択し、Advanced properties ウィンドウを展開します。Connections セクションで、先ほど作成した接続を選択します。コネクタは必要な JAR をジョブランタイムに自動的に追加します。

必要な JAR を手動で追加するには、次の手順に従ってください。
コネクタ JAR をアップロードする
最新の Spark コネクタ JAR (clickhouse-spark-runtime-3.X_2.X-0.10.X.jar) を S3 バケットにアップロードします。
S3 バケットへのアクセスを付与する
Glue ジョブがこのバケットにアクセスできることを確認します。
依存 JAR パスを設定する
Job details タブで下にスクロールし、Advanced properties のドロップダウンを展開して、Dependent JARs path に JAR のパスを入力します。

認証情報に AWS Secrets Manager を使用する
ClickHouse のユーザー名とパスワードをジョブ内にハードコードする代わりに、AWS Secrets Manager に保存し、Glue の接続またはジョブスクリプトからそのシークレットを参照します。実行時には、Glue がシークレットを取得し、そのキー・バリューのペアをコネクタの接続オプションにマージします。
シークレットを作成する
AWS Secrets Manager で、キーがコネクタのオプション名と一致するキー・バリューのペアを含む その他のシークレットタイプ のシークレットを作成します。
| Key | Value |
|---|---|
user |
ClickHouse のユーザー名 |
password |
ClickHouse のパスワード |
シークレットに追加したキーはすべてコネクタに転送されるため、コードに含めたくない場合は、host、database、そのほかのオプションもここに保存できます。
シークレットを参照する
シークレットをジョブで利用する方法は 2 つあります。
方法 1: Glue 接続にアタッチする。 Glue Studio で ClickHouse connection を作成または編集する際に、AWS secret フィールドへシークレット名を設定します。この接続を使用するジョブではシークレットが自動的に解決されるため、コードを変更する必要はありません。
方法 2: connection options で secretId を渡す。 シークレットが接続にアタッチされていない場合は、この方法を使用します。connectionName とあわせて secretId を追加します。
source = glueContext.create_dynamic_frame.from_options(
connection_type="marketplace.spark",
connection_options={
"connectionName": "<your-connection-name>",
"secretId": "clickhouse/glue/credentials",
"database": "default",
"table": "example_table"
},
transformation_ctx="clickhouse_source"
)val source = glueContext.getSource(
connectionType = "marketplace.spark",
connectionOptions = JsonOptions(Map(
"connectionName" -> "<your-connection-name>",
"secretId" -> "clickhouse/glue/credentials",
"database" -> "default",
"table" -> "example_table"
)),
transformationContext = "clickhouseSource"
)シークレット内の user キーと password キーは実行時にコネクタのオプションへマージされるため、スクリプト内でそれらを読み取る必要はありません。
例
以下の例では marketplace.spark を使用し、connectionName でコネクタを参照します。コネクタを手動でインストールした場合 (手動インストール タブ) は、代わりに connection_type="custom.spark" を使用し、className、host、http_port、user、password を options に直接指定してください。
接続自体に AWS シークレットをアタッチしている場合 (Using AWS Secrets Manager for credentials の Option 1) は、options から secretId を削除してください。Glue は接続から認証情報を自動的に取得します。
Glue Studio の visual editor では、ClickHouse コネクタをソースとしてもターゲットとしても使用できます。ClickHouse Spark Connector コンポーネントをキャンバスにドラッグし、データパイプラインに接続するだけです。

import com.amazonaws.services.glue.GlueContext
import com.amazonaws.services.glue.util.{GlueArgParser, Job, JsonOptions}
import org.apache.spark.SparkContext
import scala.collection.JavaConverters._
object ClickHouseGlueExample {
def main(sysArgs: Array[String]): Unit = {
val args = GlueArgParser.getResolvedOptions(sysArgs, Seq("JOB_NAME").toArray)
val sc = new SparkContext()
val glueContext = new GlueContext(sc)
Job.init(args("JOB_NAME"), glueContext, args.asJava)
val readOptions = JsonOptions(Map(
"connectionName" -> "<your-connection-name>",
"secretId" -> "clickhouse/glue/credentials",
"database" -> "default",
"table" -> "example_table"
))
val source = glueContext.getSource(
connectionType = "marketplace.spark",
connectionOptions = readOptions,
transformationContext = "clickhouseSource"
)
val dyf = source.getDynamicFrame()
val writeOptions = JsonOptions(Map(
"connectionName" -> "<your-connection-name>",
"secretId" -> "clickhouse/glue/credentials",
"database" -> "default",
"table" -> "target_table"
))
glueContext.getSink(
connectionType = "marketplace.spark",
connectionOptions = writeOptions
).writeDynamicFrame(dyf)
Job.commit()
}
}import sys
from awsglue.utils import getResolvedOptions
from pyspark.context import SparkContext
from awsglue.context import GlueContext
from awsglue.job import Job
args = getResolvedOptions(sys.argv, ['JOB_NAME'])
sc = SparkContext()
glueContext = GlueContext(sc)
logger = glueContext.get_logger()
job = Job(glueContext)
job.init(args['JOB_NAME'], args)
read_options = {
"connectionName": "<your-connection-name>",
"secretId": "clickhouse/glue/credentials",
"database": "default",
"table": "example_table"
}
source = glueContext.create_dynamic_frame.from_options(
connection_type="marketplace.spark",
connection_options=read_options,
transformation_ctx="clickhouse_source"
)
dyf = source
logger.info(f"Read {dyf.count()} rows from ClickHouse")
write_options = {
"connectionName": "<your-connection-name>",
"secretId": "clickhouse/glue/credentials",
"database": "default",
"table": "target_table"
}
glueContext.write_dynamic_frame.from_options(
frame=dyf,
connection_type="marketplace.spark",
connection_options=write_options,
transformation_ctx="clickhouse_sink"
)
job.commit()詳細については、Spark ドキュメント を参照してください。