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
After processing, the target contains:
Target
101
102
104
105
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
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)
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"
)
Think about it like this:
SOURCE TARGET
101 ──────────────── 101
102 ──────────────── 102
103 ─────── X -
104 ──────────────── 104
105 ──────────────── 105
The result is:
103
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"
)
For example:
SOURCE
101
102
103
104
TARGET
101
102
103
104
105
The result is:
105
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"
)
Then:
reconciliation_df = (
missing_in_target
.withColumn("issue", lit("MISSING_IN_TARGET"))
.unionByName(
extra_in_target
.withColumn("issue", lit("EXTRA_IN_TARGET"))
)
)
Now we have a single reconciliation result.
For example:
id issue
----------------------
103 MISSING_IN_TARGET
108 EXTRA_IN_TARGET
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())
Counts are useful, but they aren't enough.
Consider this:
Source count = 100,000
Target count = 100,000
It looks perfect.
But the actual records could be:
Source → 1, 2, 3, 4, 5 ...
Target → 1, 2, 3, 4, 999 ...
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"
)
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
The reconciliation output can then be written back to S3:
missing_in_target.write \
.mode("overwrite") \
.parquet(reconciliation_path)
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"
)
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)