安装
- 从 AWS Marketplace 安装 ClickHouse Glue 连接器 (推荐) 。
- 手动将 Spark Connector 的 JAR 添加到你的 Glue 作业中。
- AWS Marketplace
- 手动安装
-
订阅连接器
要在你的账户中使用该连接器,请先在 AWS Marketplace 订阅 ClickHouse AWS Glue Connector。 -
授予所需权限
确保你的 Glue 作业所使用的 IAM role 具备所需权限,具体请参见最低权限指南。 -
激活连接器并创建连接
你可以点击此链接直接激活连接器并创建连接。该链接会打开 Glue 连接创建页面,并预先填好关键字段。为连接命名后,点击 create 即可 (此阶段无需提供 ClickHouse 连接信息) 。 -
在 Glue 作业中使用
在你的 Glue 作业中,选择Job details选项卡,然后展开Advanced properties窗口。在Connections部分中,选择你刚刚创建的连接。该连接器会自动将所需 JAR 注入作业运行时。
Glue 连接器中使用的 JAR 是为
Spark 3.3、Scala 2 和 Python 3 构建的。配置 Glue 作业时,请确保选择这些版本。如需手动添加所需的 JAR,请按以下步骤操作:
- 将以下 JAR 上传到一个 S3 bucket:
clickhouse-jdbc-0.6.X-all.jar和clickhouse-spark-runtime-3.X_2.X-0.8.X.jar。 - 确保 Glue 作业可以访问该 bucket。
- 在
Job details选项卡下,向下滚动并展开Advanced properties下拉菜单,然后在Dependent JARs path中填写 JAR 路径:
示例
- Scala
- Python
import com.amazonaws.services.glue.GlueContext
import com.amazonaws.services.glue.util.GlueArgParser
import com.amazonaws.services.glue.util.Job
import com.clickhouseScala.Native.NativeSparkRead.spark
import org.apache.spark.sql.SparkSession
import scala.collection.JavaConverters._
import org.apache.spark.sql.types._
import org.apache.spark.sql.functions._
object ClickHouseGlueExample {
def main(sysArgs: Array[String]) {
val args = GlueArgParser.getResolvedOptions(sysArgs, Seq("JOB_NAME").toArray)
val sparkSession: SparkSession = SparkSession.builder
.config("spark.sql.catalog.clickhouse", "com.clickhouse.spark.ClickHouseCatalog")
.config("spark.sql.catalog.clickhouse.host", "<your-clickhouse-host>")
.config("spark.sql.catalog.clickhouse.protocol", "https")
.config("spark.sql.catalog.clickhouse.http_port", "<your-clickhouse-port>")
.config("spark.sql.catalog.clickhouse.user", "default")
.config("spark.sql.catalog.clickhouse.password", "<your-password>")
.config("spark.sql.catalog.clickhouse.database", "default")
// for ClickHouse cloud
.config("spark.sql.catalog.clickhouse.option.ssl", "true")
.config("spark.sql.catalog.clickhouse.option.ssl_mode", "NONE")
.getOrCreate
val glueContext = new GlueContext(sparkSession.sparkContext)
Job.init(args("JOB_NAME"), glueContext, args.asJava)
import sparkSession.implicits._
val url = "s3://{path_to_cell_tower_data}/cell_towers.csv.gz"
val schema = StructType(Seq(
StructField("radio", StringType, nullable = false),
StructField("mcc", IntegerType, nullable = false),
StructField("net", IntegerType, nullable = false),
StructField("area", IntegerType, nullable = false),
StructField("cell", LongType, nullable = false),
StructField("unit", IntegerType, nullable = false),
StructField("lon", DoubleType, nullable = false),
StructField("lat", DoubleType, nullable = false),
StructField("range", IntegerType, nullable = false),
StructField("samples", IntegerType, nullable = false),
StructField("changeable", IntegerType, nullable = false),
StructField("created", TimestampType, nullable = false),
StructField("updated", TimestampType, nullable = false),
StructField("averageSignal", IntegerType, nullable = false)
))
val df = sparkSession.read
.option("header", "true")
.schema(schema)
.csv(url)
// Write to ClickHouse
df.writeTo("clickhouse.default.cell_towers").append()
// Read from ClickHouse
val dfRead = spark.sql("select * from clickhouse.default.cell_towers")
Job.commit()
}
}
import sys
from awsglue.transforms import *
from awsglue.utils import getResolvedOptions
from pyspark.context import SparkContext
from awsglue.context import GlueContext
from awsglue.job import Job
from pyspark.sql import Row
## @params: [JOB_NAME]
args = getResolvedOptions(sys.argv, ['JOB_NAME'])
sc = SparkContext()
glueContext = GlueContext(sc)
logger = glueContext.get_logger()
spark = glueContext.spark_session
job = Job(glueContext)
job.init(args['JOB_NAME'], args)
spark.conf.set("spark.sql.catalog.clickhouse", "com.clickhouse.spark.ClickHouseCatalog")
spark.conf.set("spark.sql.catalog.clickhouse.host", "<your-clickhouse-host>")
spark.conf.set("spark.sql.catalog.clickhouse.protocol", "https")
spark.conf.set("spark.sql.catalog.clickhouse.http_port", "<your-clickhouse-port>")
spark.conf.set("spark.sql.catalog.clickhouse.user", "default")
spark.conf.set("spark.sql.catalog.clickhouse.password", "<your-password>")
spark.conf.set("spark.sql.catalog.clickhouse.database", "default")
spark.conf.set("spark.clickhouse.write.format", "json")
spark.conf.set("spark.clickhouse.read.format", "arrow")
# for ClickHouse cloud
spark.conf.set("spark.sql.catalog.clickhouse.option.ssl", "true")
spark.conf.set("spark.sql.catalog.clickhouse.option.ssl_mode", "NONE")
# Create DataFrame
data = [Row(id=11, name="John"), Row(id=12, name="Doe")]
df = spark.createDataFrame(data)
# Write DataFrame to ClickHouse
df.writeTo("clickhouse.default.example_table").append()
# Read DataFrame from ClickHouse
df_read = spark.sql("select * from clickhouse.default.example_table")
logger.info(str(df.take(10)))
job.commit()