新手求助:基于Scala+Deequ计算S3 CSV指标并写入Glue表
Answer
Hi there! As someone new to Scala and Deequ, let’s break this down step by step with practical code and explanations to help you achieve your goal.
Prerequisites
First, make sure your project includes the necessary dependencies in build.sbt:
libraryDependencies ++= Seq( "org.apache.spark" %% "spark-sql" % "3.3.2" % Provided, "com.amazon.deequ" %% "deequ" % "2.0.0-spark-3.3", "software.amazon.awssdk" % "s3" % "2.20.100", "org.apache.hadoop" % "hadoop-aws" % "3.3.4", "com.amazonaws" % "aws-java-sdk-glue" % "1.12.471" )
Note: Match the Spark and Deequ versions (Deequ’s version format follows [deequ-version]-spark-[spark-version] to ensure compatibility).
Step-by-Step Example Code
Here’s a complete Scala script that reads your CSV from S3, runs Deequ quality checks, and loads the resulting metrics to a Glue table:
import org.apache.spark.sql.SparkSession import com.amazon.deequ.VerificationSuite import com.amazon.deequ.checks.{Check, CheckLevel, CheckStatus} import org.apache.spark.sql.types.{StringType, StructField, StructType} object DeequS3ToGlue { def main(args: Array[String]): Unit = { // 1. Initialize Spark Session with Glue Catalog support val spark = SparkSession.builder() .appName("DeequMetricsToGlue") .config("spark.sql.catalogImplementation", "hive") .config("hive.metastore.client.factory.class", "com.amazonaws.glue.catalog.metastore.AWSGlueDataCatalogHiveClientFactory") .getOrCreate() import spark.implicits._ // 2. Read CSV from S3 (update path and schema to match your data) val inputSchema = StructType(Seq( StructField("user_id", StringType, nullable = false), StructField("email", StringType, nullable = true), StructField("country", StringType, nullable = true) )) val sourceDF = spark.read .schema(inputSchema) .option("header", "true") .csv("s3://your-source-bucket/path/to/data.csv") // 3. Define Deequ checks for Completeness and CountDistinct val verificationResult = VerificationSuite() .onData(sourceDF) .addCheck( Check(CheckLevel.Error, "Core Data Quality Checks") // Completeness: Percentage of non-null values in a column .isComplete("user_id") // Require 100% completeness for user_id .isComplete("email", 0.95) // Allow up to 5% nulls for email // CountDistinct: Number of unique values in a column .hasDistinctCount("user_id", _ >= sourceDF.count()) // Ensure no duplicate user_ids .hasDistinctCount("country", _ >= 10) // Require at least 10 unique countries ) .run() // 4. Convert metrics to a structured DataFrame for Glue val metricsDF = verificationResult.metrics.map { case (_, metric) => ( metric.entity, // e.g., "Column" metric.instance, // e.g., "user_id" metric.name, // e.g., "Completeness" metric.value.getOrElse(-1.0) // Numeric value of the metric ) }.toSeq.toDF("entity_type", "column_name", "metric_name", "metric_value") // 5. Write metrics to Glue Table (update database and table names) metricsDF.write .mode("overwrite") // Use "append" if you want to accumulate metrics over time .format("parquet") .option("path", "s3://your-metrics-bucket/path/to/glue-table-storage/") .saveAsTable("your_glue_database.data_quality_metrics") // Print check status for quick debugging if (verificationResult.status == CheckStatus.Success) { println("✅ All data quality checks passed!") } else { println("❌ Some data quality checks failed — check metrics for details.") } spark.stop() } }
Key Deequ Concepts Explained
Let’s clarify the core parts you’ll work with:
- VerificationSuite: The main entry point to execute data quality checks. It takes your input DataFrame and runs the defined rules.
- Check: A container for quality rules. You set a severity level (
Error/Warning) and add constraints:isComplete(column, minThreshold): Validates that the percentage of non-null values incolumnmeets or exceedsminThreshold(defaults to 1.0 = 100%).hasDistinctCount(column, constraint): Checks the number of unique values incolumnagainst a condition (e.g.,_ >= 10).
- VerificationResult: Stores the pass/fail status of checks and all computed metrics, which we convert to a structured DataFrame for loading into Glue.
Loading Metrics to Glue Table Notes
- IAM Permissions: Ensure your Spark job’s IAM role has:
s3:GetObjectaccess to the source S3 bucket.s3:PutObjectaccess to the metrics storage bucket.glue:CreateTable/glue:UpdateTablepermissions for your target Glue database.
- Table Registration: The
saveAsTablemethod automatically creates the Glue table if it doesn’t exist. If it does exist,mode("overwrite")replaces the data (useappendto build a historical metrics log). - Storage Format: We use Parquet because it’s columnar, efficient for analytics, and fully supported by Glue and Athena.
Troubleshooting Tips
- S3 Access Errors: Double-check your IAM role’s bucket policies and make sure the Spark job is running in a region where the buckets exist.
- Glue Table Not Visible: Confirm the Spark session is configured with the Glue catalog factory (the
hive.metastore.client.factory.classsetting is critical here). - Dependency Conflicts: If running on EMR or Glue Jobs, mark Spark dependencies as
Providedsince these environments include Spark out of the box.
内容的提问来源于stack exchange,提问作者marie20
相关产品推荐
相关产品推荐

