AWS Glue中Scala多S3文件关联及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

