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

新手求助:基于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 in column meets or exceeds minThreshold (defaults to 1.0 = 100%).
    • hasDistinctCount(column, constraint): Checks the number of unique values in column against 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:GetObject access to the source S3 bucket.
    • s3:PutObject access to the metrics storage bucket.
    • glue:CreateTable/glue:UpdateTable permissions for your target Glue database.
  • Table Registration: The saveAsTable method automatically creates the Glue table if it doesn’t exist. If it does exist, mode("overwrite") replaces the data (use append to 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.class setting is critical here).
  • Dependency Conflicts: If running on EMR or Glue Jobs, mark Spark dependencies as Provided since these environments include Spark out of the box.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.14 07:41:28