DEV Community

Cover image for Why Data Reconciliation Matters in AWS Glue: Finding Missing Records with Anti Joins
Rahul R
Rahul R

Posted on

Why Data Reconciliation Matters in AWS Glue: Finding Missing Records with Anti Joins

A pipeline can be green and still have a data problem.

This is one of the things I learned while working with data pipelines.

A Glue job can complete successfully, files can be written to S3, and there may be no obvious error in the logs. But that doesn't necessarily mean that every record made it from the source to the target.

That's where data reconciliation comes in.

What is data reconciliation?

In simple terms, reconciliation means comparing two datasets to check whether the data moved as expected.

Imagine we have data coming from a source system:

Source
101
102
103
104
105
Enter fullscreen mode Exit fullscreen mode

After processing, the target contains:

Target
101
102
104
105
Enter fullscreen mode Exit fullscreen mode

The Glue job might still be marked as successful.

But record 103 is missing.

A reconciliation check helps us catch this instead of assuming that a successful pipeline means successful data movement.


Why do we need reconciliation?

There are several situations where records can disappear or become different during a pipeline:

  • Filtering conditions
  • Incorrect joins
  • Transformation logic
  • Duplicate handling
  • Incremental-load issues
  • Source/target timing differences
  • Failed or partially processed batches
  • Incorrect configuration

Checking only the job status doesn't tell us whether the actual data is complete.

That's why I think of reconciliation as a second layer of validation:

Source
   ↓
AWS Glue
   ↓
Transformations
   ↓
Target
   ↓
Reconciliation
   ↓
Missing / Extra / Mismatched Records
Enter fullscreen mode Exit fullscreen mode

Using AWS Glue for reconciliation

AWS Glue works with Spark, so we can use Spark DataFrame operations to compare datasets.

For example, let's assume both datasets have a business key called id.

source_df = spark.read.parquet(source_path)
target_df = spark.read.parquet(target_path)
Enter fullscreen mode Exit fullscreen mode

Now we can compare them.

One of the most useful tools for this is an anti join.


1. Left Anti Join

A left anti join returns records that exist in the left DataFrame but don't exist in the right DataFrame.

missing_in_target = source_df.join(
    target_df,
    on="id",
    how="left_anti"
)
Enter fullscreen mode Exit fullscreen mode

Think about it like this:

SOURCE                TARGET

101  ──────────────── 101
102  ──────────────── 102
103  ─────── X        -
104  ──────────────── 104
105  ──────────────── 105
Enter fullscreen mode Exit fullscreen mode

The result is:

103
Enter fullscreen mode Exit fullscreen mode

So this answers:

Which records exist in my source but are missing from my target?

This is probably the most common reconciliation check.


2. Checking the Other Direction

Sometimes we also want to know whether the target contains records that don't exist in the source.

We can simply reverse the DataFrames:

extra_in_target = target_df.join(
    source_df,
    on="id",
    how="left_anti"
)
Enter fullscreen mode Exit fullscreen mode

For example:

SOURCE

101
102
103
104

TARGET

101
102
103
104
105
Enter fullscreen mode Exit fullscreen mode

The result is:

105
Enter fullscreen mode Exit fullscreen mode

This tells us:

Which records exist in the target but not in the source?

This is useful for finding unexpected or extra records.


3. What about a "Full Anti Join"?

Spark doesn't have a join type literally called full_anti.

Instead, we can perform both directions and combine the results.

missing_in_target = source_df.join(
    target_df,
    on="id",
    how="left_anti"
)

extra_in_target = target_df.join(
    source_df,
    on="id",
    how="left_anti"
)
Enter fullscreen mode Exit fullscreen mode

Then:

reconciliation_df = (
    missing_in_target
    .withColumn("issue", lit("MISSING_IN_TARGET"))
    .unionByName(
        extra_in_target
        .withColumn("issue", lit("EXTRA_IN_TARGET"))
    )
)
Enter fullscreen mode Exit fullscreen mode

Now we have a single reconciliation result.

For example:

id     issue
----------------------
103    MISSING_IN_TARGET
108    EXTRA_IN_TARGET
Enter fullscreen mode Exit fullscreen mode

This gives us a more complete picture of the difference between the two datasets.


4. Don't stop at record counts

A common first check is:

print(source_df.count())
print(target_df.count())
Enter fullscreen mode Exit fullscreen mode

Counts are useful, but they aren't enough.

Consider this:

Source count = 100,000
Target count = 100,000
Enter fullscreen mode Exit fullscreen mode

It looks perfect.

But the actual records could be:

Source → 1, 2, 3, 4, 5 ...

Target → 1, 2, 3, 4, 999 ...
Enter fullscreen mode Exit fullscreen mode

The counts match, but the data doesn't.

That's why comparing the actual business keys is much more useful.

missing = source_df.join(
    target_df,
    "id",
    "left_anti"
)
Enter fullscreen mode Exit fullscreen mode

5. Reconciliation in a Glue pipeline

A simple architecture could look like this:

Source System
      |
      ↓
   S3 Raw
      |
      ↓
   AWS Glue
      |
      ↓
 Transformations
      |
      ↓
  S3 Target
      |
      ↓
 Reconciliation
      |
      ├── Missing records
      ├── Extra records
      └── Reconciliation report
Enter fullscreen mode Exit fullscreen mode

The reconciliation output can then be written back to S3:

missing_in_target.write \
    .mode("overwrite") \
    .parquet(reconciliation_path)
Enter fullscreen mode Exit fullscreen mode

Depending on the requirement, the reconciliation result could also be used to trigger an alert or create a report for further investigation.


A small but important detail

The join key matters.

If id isn't unique, simply comparing IDs may not be enough.

Sometimes reconciliation needs multiple columns:

join_columns = ["customer_id", "order_id", "transaction_date"]

missing = source_df.join(
    target_df,
    on=join_columns,
    how="left_anti"
)
Enter fullscreen mode Exit fullscreen mode

The right reconciliation key depends on the business meaning of the data.


Final thought

For me, reconciliation is not just another validation step.

It's a way of answering a simple question:

"Did the data I expected to move actually move?"

A successful Glue job tells us that the processing completed.

A reconciliation check tells us whether the data itself looks complete.

And that's an important difference.

Anti joins are especially useful here because they let us directly find the records that exist on one side but not the other, instead of only comparing counts.

Once you start using reconciliation as part of the pipeline, a "SUCCESS" status becomes more meaningful because you're checking both the process and the data.

Top comments (0)