You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

AWS Glue中Scala多S3文件关联及MySQL增改示例脚本请求

Scala Example for Joining S3 Files in AWS Glue and Upserting to MySQL

Got it, here's a practical Scala example tailored for your AWS Glue workflow: joining multiple S3 datasets and performing insert/update (upsert) operations on a MySQL database. I’ll walk you through the full script and explain key parts so you can adapt it to your specific schema.

Prerequisites

First, make sure you have these setup:

  • An AWS Glue IAM role with:
    • Read access to your target S3 buckets
    • Permissions to connect to your MySQL instance (either via VPC access or public endpoint, depending on your setup)
  • MySQL JDBC driver available in Glue (you can use the default one or upload a custom version to S3 and reference it in your job)
  • Your S3 files are in a structured format like Parquet, CSV, or JSON (we’ll use Parquet in this example, but you can adjust for other formats)

Full Scala Script

import com.amazonaws.services.glue.GlueContext
import com.amazonaws.services.glue.util.GlueArgParser
import com.amazonaws.services.glue.util.Job
import com.amazonaws.services.glue.util.JsonOptions
import org.apache.spark.SparkContext
import org.apache.spark.sql.{DataFrame, SaveMode}
import org.apache.spark.sql.functions._

object GlueS3JoinAndUpsertToMySQL {
  def main(sysArgs: Array[String]) {
    val sc: SparkContext = new SparkContext()
    val glueContext: GlueContext = new GlueContext(sc)
    val spark = glueContext.getSparkSession
    
    // Parse job arguments (optional, but useful for dynamic paths)
    val args = GlueArgParser.getResolvedOptions(sysArgs, Seq("JOB_NAME", "s3_path_orders", "s3_path_customers", "mysql_url", "mysql_table", "mysql_user", "mysql_password").toArray)
    Job.init(args("JOB_NAME"), glueContext, args.asJava)

    // --------------------------
    // Step 1: Read S3 datasets
    // --------------------------
    // Read orders data from S3 (Parquet format)
    val ordersDF: DataFrame = glueContext.getSourceWithFormat(
      formatOptions = JsonOptions(Map("compression" -> "snappy")),
      connectionType = "s3",
      format = "parquet",
      options = JsonOptions(Map("path" -> args("s3_path_orders")))
    ).getDynamicFrame().toDF()

    // Read customers data from S3 (Parquet format)
    val customersDF: DataFrame = glueContext.getSourceWithFormat(
      formatOptions = JsonOptions(Map("compression" -> "snappy")),
      connectionType = "s3",
      format = "parquet",
      options = JsonOptions(Map("path" -> args("s3_path_customers")))
    ).getDynamicFrame().toDF()

    // --------------------------
    // Step 2: Join the datasets
    // --------------------------
    // Join orders with customers on customer_id (adjust join key and type as needed)
    val joinedDF: DataFrame = ordersDF.join(
      customersDF,
      ordersDF("customer_id") === customersDF("customer_id"),
      "inner" // Use "left", "right", or "full" depending on your needs
    )
    .select(
      ordersDF("order_id"),
      ordersDF("order_date"),
      ordersDF("total_amount"),
      customersDF("customer_id"),
      customersDF("customer_name"),
      customersDF("email")
    )
    // Add a last_updated timestamp for tracking
    .withColumn("last_updated", current_timestamp())

    // --------------------------
    // Step 3: Upsert to MySQL
    // --------------------------
    // Define MySQL connection properties
    val mysqlProps = Map(
      "user" -> args("mysql_user"),
      "password" -> args("mysql_password"),
      "driver" -> "com.mysql.cj.jdbc.Driver"
    )

    // Set MySQL connection for Spark
    spark.sqlContext.setConf("spark.sql.catalog.mysql_catalog", "org.apache.spark.sql.jdbc.JdbcCatalog")
    spark.sqlContext.setConf("spark.sql.catalog.mysql_catalog.url", args("mysql_url"))
    spark.sqlContext.setConf("spark.sql.catalog.mysql_catalog.user", args("mysql_user"))
    spark.sqlContext.setConf("spark.sql.catalog.mysql_catalog.password", args("mysql_password"))

    // Use foreachBatch to handle upserts (since Spark's JDBC doesn't support native upsert)
    joinedDF.writeStream.foreachBatch { (batchDF: DataFrame, batchId: Long) =>
      // Create a temporary view for the batch data
      batchDF.createOrReplaceTempView("temp_joined_data")

      // Execute upsert using MySQL's INSERT ... ON DUPLICATE KEY UPDATE
      spark.sql(
        s"""
           |INSERT INTO mysql_catalog.${args("mysql_table")} 
           |(order_id, order_date, total_amount, customer_id, customer_name, email, last_updated)
           |SELECT order_id, order_date, total_amount, customer_id, customer_name, email, last_updated
           |FROM temp_joined_data
           |ON DUPLICATE KEY UPDATE
           |total_amount = VALUES(total_amount),
           |customer_name = VALUES(customer_name),
           |email = VALUES(email),
           |last_updated = VALUES(last_updated)
         """.stripMargin
      )
    }
    .option("checkpointLocation", "s3://your-checkpoint-bucket/glue-checkpoints/") // Required for streaming
    .trigger(once = true) // Run once for batch processing; adjust for streaming
    .start()
    .awaitTermination()

    Job.commit()
  }
}

Key Explanations

1. Reading S3 Data

We use Glue's getSourceWithFormat to read structured data from S3. For CSV files, you’d adjust the format to "csv" and add formatOptions like Map("header" -> "true", "delimiter" -> ",").

2. Joining Datasets

The example uses an inner join on customer_id, but you can change the join type (left, right, full) based on your data requirements. We also select only the columns we need to avoid unnecessary data transfer.

3. Upsert Logic

Spark's native JDBC writer doesn’t support upserts directly, so we use foreachBatch to process each batch of data. We create a temporary view and execute a MySQL-specific INSERT ... ON DUPLICATE KEY UPDATE statement. Make sure your MySQL table has a primary key (like order_id) so the database can identify which rows to update.

Important Notes

  • Permissions: Ensure your Glue role has the AmazonS3FullAccess (or restricted to specific buckets) and permissions to connect to your MySQL instance (if using VPC, the role needs access to the VPC resources).
  • Data Types: Verify that the data types in your S3 files match the MySQL table schema to avoid conversion errors.
  • Performance: For large datasets, consider:
    • Reading S3 data in partitions to reduce data load
    • Broadcasting small tables (use broadcast(customersDF) if customers is smaller than orders)
    • Using Glue's DynamicFrame transformations if you need schema evolution handling
  • Checkpoint Location: For batch processing with trigger(once=true), the checkpoint location is still required but will only be used once. For continuous streaming, it tracks processed data to avoid duplicates.

内容的提问来源于stack exchange,提问作者Ramesh Daddala

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.21 07:50:40