Apache Hudi в AWS Glue
Вы когда-нибудь задавались вопросом, как писать таблицы Hudi (Scala) в AWS Glue?
Если да, то Вы попали по адресу.
Необходимые условия
- Создайте базу данных Glue под названием hudi_db в меню Databases в Data Catalogue в консоли Glue Console
Алгоритм действий
Настройка задания
- В консоли Glue выберите Задания ETL, затем выберите Script Editor.
- Теперь на вкладке выше выберите Job details, в Language выберите Scala.
- Не бойтесь вносить любые изменения в инфраструктуру по мере необходимости.
-
Нажмите на Advanced properties, перейдите к Job parameters и поочередно добавьте все параметры, указанные ниже. Измените эти переменные по своему усмотрению.
-
--S3_OUTPUT_PATHass3://hudi-spark-quickstart/write-path/ -
--classasGlueApp -
--confasspark.serializer=org.apache.spark.serializer.KryoSerializer --conf spark.sql.hive.convertMetastoreParquet=false -
--datalake-formatsashudi
-
Примечание: В данном случае я использую версию Hudi 0.12.0, которая по умолчанию поставляется вместе с Glue 4.0. Если Вы хотите использовать другую версию Hudi, Вам придется добавить jar в путь класса, добавив еще одно свойство --extra-jars и указав на S3 путь к JAR-файлу Hudi.
Переходим к самому интересному.
Написание скрипта
Перейдите на вкладку Script и добавьте код на языке Scala, приведенный ниже.
import com.amazonaws.services.glue.{GlueContext, DynamicFrame}
import com.amazonaws.services.glue.util.GlueArgParser
import org.apache.spark.SparkContext
import org.apache.spark.sql.{SQLContext, SparkSession}
import org.apache.spark.sql.types._
import org.apache.hudi.DataSourceWriteOptions._
import org.apache.hudi.config.HoodieWriteConfig._
import com.amazonaws.services.glue.log.GlueLogger
Добавьте код, специфичный для Glue, т.е. для анализа параметров задания и создания glueContext.
object GlueApp {
def main(sysArgs: Array[String]) {
val args = GlueArgParser.getResolvedOptions(sysArgs, Seq("JOB_NAME", "S3_OUTPUT_PATH").toArray)
val spark: SparkSession = SparkSession.builder().appName("AWS Glue Hudi Job").getOrCreate()
val glueContext: GlueContext = new GlueContext(spark.sparkContext)
val logger = new GlueLogger()Подготовка данных
import spark.implicits._
val tableName = "trips"
val recordKeyColumn = "uuid" val precombineKeyColumn = "ts" val partitionKeyColumn = "city" val s3OutputPath = args("S3_OUTPUT_PATH") val glueDbName = "hudi_db" val writePath = s"$s3OutputPath/$tableName" val columns = Seq("ts","uuid","rider","driver","fare","city") val data = Seq((1695159649087L,"334e26e9-8355-45cc-97c6-c31daf0df330","rider-A","driver-K",19.10,"san_francisco"), (1695091554788L,"e96c4396-3fad-413a-a942-4cb36106d721","rider-C","driver-M",27.70 ,"san_francisco"), (1695046462179L,"9909a8b1-2d15-4d3d-8ec9-efc48c536a00","rider-D","driver-L",33.90 ,"san_francisco"), (1695516137016L,"e3cf430c-889d-4015-bc98-59bdce1e530c","rider-F","driver-P",34.15,"sao_paulo" ), (1695115999911L,"c8abbe79-8d89-47ea-b4ce-4d224bae5bfa","rider-J","driver-T",17.85,"chennai"));
Добавьте параметры, необходимые Hudi для записи таблицы и синхронизации ее с базой данных Glue Database.
val hudiOptions = Map[String, String](
"hoodie.table.name" -> tableName,
"hoodie.datasource.write.recordkey.field" -> recordKeyColumn,
"hoodie.datasource.write.precombine.field" -> precombineKeyColumn,
"hoodie.datasource.write.partitionpath.field" -> partitionKeyColumn,
"hoodie.datasource.write.hive_style_partitioning" -> "true",
"hoodie.datasource.write.storage.type" -> "COPY_ON_WRITE",
"hoodie.datasource.write.operation" -> "upsert",
"hoodie.datasource.hive_sync.enable" -> "true",
"hoodie.datasource.hive_sync.database" -> glueDbName,
"hoodie.datasource.hive_sync.table" -> tableName,
"hoodie.datasource.hive_sync.partition_fields" -> partitionKeyColumn,
"hoodie.datasource.hive_sync.partition_extractor_class" -> "org.apache.hudi.hive.MultiPartKeysValueExtractor",
"hoodie.datasource.hive_sync.use_jdbc" -> "false",
"hoodie.datasource.hive_sync.mode" -> "hms",
"path" -> writePath
)
Наконец, создайте датафрейм и занесите его в S3.
var inserts = spark.createDataFrame(data).toDF(columns:_*)
inserts.write
.format("hudi") .options(hudiOptions) .mode("overwrite") .save() logger.info("Data successfully written to S3 using Hudi") } }Отправка запросов
Теперь, когда мы записали таблицу в S3, мы можем запросить ее из Athena.
SELECT * FROM "hudi_db"."trips" limit 10;
Вот и все!




