K
Kiran
Guest
In December 2020, AWS made S3 strongly read-after-write consistent. A PUT shows up in the next LIST. A GET after an overwrite returns the new bytes.
Some months ago a Glue job of mine started dropping records from a 6 MB CSV, and this guarantee sent me looking in the wrong direction.
The job read one object from an S3 prefix. Same key every run,
I saved one of the short results. On a rerun the missing records were back.
I assumed Spark had read the object while it was being written. It had not.
I investigated the runs as far as I could, but the evidence was not enough to establish a single root cause. Two mechanisms remained consistent with what I had observed.
S3's guarantee is per object and per operation. Once a PUT completes, a subsequent GET sees the new object. Once an object is created or deleted, LIST reflects it.
That does not make a sequence of operations a snapshot.
A job can list a path, record what it found, do other work, and open the object later. If somebody overwrites the key in between, both S3 operations can be strongly consistent and still refer to different versions of the object.
I had been treating those as the same guarantee.
A partial PUT was also the wrong explanation. S3 does not expose the first part of an in-progress
So the fact that Spark returned partial data did not mean S3 had returned a partial object.
Spark lists files when it resolves the input and records metadata about them, including their lengths. That length later becomes part of the file split used by the task.
Roughly:
The listing and the actual read do not have to happen at the same time.
Say Spark lists
The GET sees the 5.8 MB version, as it should. But Spark planned the work using 4.1 MB. It reads to the boundary it already has. There is no task for the remaining 1.7 MB.
No S3 consistency violation is required.
This also explains why the format mattered. A Parquet file has enough internal structure that truncation tends to fail loudly. CSV does not. Stop at a line boundary and everything read so far is still valid CSV. Even a final fragment can disappear depending on the parser mode.
The sizes above are only an example. I could not establish the sizes of the object versions involved in the bad runs, so this mechanism remains an explanation rather than a confirmed root cause. If the relevant versions were all the same byte length, it cannot explain what happened.
Two CSVs containing the same logical data do not necessarily have the same byte length. Row order, line endings, quoting, timestamp formatting, a trailing newline, or a BOM can all change the size. An intermediate output that is later replaced by the complete output would do the same.
This code makes the race look smaller than it is:
The input metadata can be resolved when the DataFrame is constructed. The bytes are needed when an action runs. Anything that happens between those points increases the window in which the key can change.
That detail cost me some time because repeated actions against the DataFrame gave me the same short result. It looked deterministic, so I kept looking at the contents of the CSV.
I checked other explanations as well.
A stale S3 read could give me an older complete version, but it would not explain reading only part of the new one. Glue bookmarks did not fit either: the records came back when I reran the job. And there is no client-side switch for old S3 consistency behavior; the consistency guarantee is service-side.
There is a similar problem with larger text files. Once a CSV is split into multiple ranges, separate tasks can read separate parts of the object. The line reader handles the boundary between splits assuming both readers see the same bytes. If the object changes underneath them, that assumption is no longer safe.
That was not my 6 MB case, but it comes from the same assumption that the object remains unchanged while the read is in progress.
The other mechanism I could not rule out is a retry across versions.
Suppose a GET is interrupted and the client retries after the object has been overwritten. If the retry is not conditional on the ETag or version it originally opened, it can continue against different bytes using an offset derived from the first request.
Whether that is possible in practice depends on the connector and version. S3A has change-detection support that can make this fail instead of silently continuing. A Glue job may also be reading through a different connector, in which case S3A settings are irrelevant.
The evidence I collected was not enough to determine whether this happened.
One useful check, where S3 version history is available, is:
If the versions have different sizes, the last record in a short output can be located in the complete object. If its end offset is close to the size of an earlier version, that would support the recorded-length explanation.
If the relevant versions are all the same size, that explanation is out.
I did not have enough historical evidence from the affected runs to make that determination conclusively.
I changed the way the input was written.
Instead of overwriting the same location on every run, the upstream process creates a new partition with an automatically generated ID. Each run therefore gets its own location rather than replacing the object that Spark may already have listed.
Conceptually, the layout changed from something like:
to:
Once Spark starts reading a partition, that partition is not being replaced by the next run. The next run writes to a different partition.
After that change, the short reads stopped.
There are other ways to avoid the same race. A staging copy, a completion protocol, or a table format with committed snapshots can all provide a stable input. But those are alternatives, not what I did here. For this job, giving every run its own partition was enough.
The part that initially confused me was the S3 consistency guarantee. Every S3 request can be strongly consistent and this failure can still happen, because Spark does not perform the listing and all subsequent reads as one atomic operation.
I collected all the evidence I could from the affected runs. It narrowed the problem considerably, but it did not distinguish conclusively between the two remaining mechanisms.
What I can say from the result is simpler: once the runs stopped sharing a mutable input location, the problem stopped.
Some months ago a Glue job of mine started dropping records from a 6 MB CSV, and this guarantee sent me looking in the wrong direction.
The job read one object from an S3 prefix. Same key every run,
filename.csv, overwritten in place by an upstream process. Most runs were fine. Some came back short. There was no exception. The DataFrame just had fewer rows than the file.I saved one of the short results. On a rerun the missing records were back.
I assumed Spark had read the object while it was being written. It had not.
I investigated the runs as far as I could, but the evidence was not enough to establish a single root cause. Two mechanisms remained consistent with what I had observed.
S3 consistency was not the problem
S3's guarantee is per object and per operation. Once a PUT completes, a subsequent GET sees the new object. Once an object is created or deleted, LIST reflects it.
That does not make a sequence of operations a snapshot.
A job can list a path, record what it found, do other work, and open the object later. If somebody overwrites the key in between, both S3 operations can be strongly consistent and still refer to different versions of the object.
I had been treating those as the same guarantee.
A partial PUT was also the wrong explanation. S3 does not expose the first part of an in-progress
PutObject. CompleteMultipartUpload has the same property from the reader's point of view: until it completes, the new object is not available as a partially assembled object.So the fact that Spark returned partial data did not mean S3 had returned a partial object.
The file length is recorded before the read
Spark lists files when it resolves the input and records metadata about them, including their lengths. That length later becomes part of the file split used by the task.
Roughly:
Code:
PartitionedFile(path, start = 0, length = <recorded length>)
The listing and the actual read do not have to happen at the same time.
Say Spark lists
filename.csv and sees a 4.1 MB object. Before the task opens it, the upstream job overwrites the key with a 5.8 MB version.The GET sees the 5.8 MB version, as it should. But Spark planned the work using 4.1 MB. It reads to the boundary it already has. There is no task for the remaining 1.7 MB.
No S3 consistency violation is required.
This also explains why the format mattered. A Parquet file has enough internal structure that truncation tends to fail loudly. CSV does not. Stop at a line boundary and everything read so far is still valid CSV. Even a final fragment can disappear depending on the parser mode.
The sizes above are only an example. I could not establish the sizes of the object versions involved in the bad runs, so this mechanism remains an explanation rather than a confirmed root cause. If the relevant versions were all the same byte length, it cannot explain what happened.
Two CSVs containing the same logical data do not necessarily have the same byte length. Row order, line endings, quoting, timestamp formatting, a trailing newline, or a BOM can all change the size. An intermediate output that is later replaced by the complete output would do the same.
The window can be surprisingly large
This code makes the race look smaller than it is:
Code:
df = spark.read.csv("s3://bucket/path/filename.csv")
# other work
df.count()
The input metadata can be resolved when the DataFrame is constructed. The bytes are needed when an action runs. Anything that happens between those points increases the window in which the key can change.
That detail cost me some time because repeated actions against the DataFrame gave me the same short result. It looked deterministic, so I kept looking at the contents of the CSV.
I checked other explanations as well.
A stale S3 read could give me an older complete version, but it would not explain reading only part of the new one. Glue bookmarks did not fit either: the records came back when I reran the job. And there is no client-side switch for old S3 consistency behavior; the consistency guarantee is service-side.
There is a similar problem with larger text files. Once a CSV is split into multiple ranges, separate tasks can read separate parts of the object. The line reader handles the boundary between splits assuming both readers see the same bytes. If the object changes underneath them, that assumption is no longer safe.
That was not my 6 MB case, but it comes from the same assumption that the object remains unchanged while the read is in progress.
The other possibility
The other mechanism I could not rule out is a retry across versions.
Suppose a GET is interrupted and the client retries after the object has been overwritten. If the retry is not conditional on the ETag or version it originally opened, it can continue against different bytes using an offset derived from the first request.
Whether that is possible in practice depends on the connector and version. S3A has change-detection support that can make this fail instead of silently continuing. A Glue job may also be reading through a different connector, in which case S3A settings are irrelevant.
The evidence I collected was not enough to determine whether this happened.
One useful check, where S3 version history is available, is:
Code:
aws s3api list-object-versions \
--bucket BUCKET \
--prefix path/filename.csv \
--query 'Versions[].{Size:Size,Modified:LastModified,Id:VersionId}'
If the versions have different sizes, the last record in a short output can be located in the complete object. If its end offset is close to the size of an earlier version, that would support the recorded-length explanation.
If the relevant versions are all the same size, that explanation is out.
I did not have enough historical evidence from the affected runs to make that determination conclusively.
What fixed it
I changed the way the input was written.
Instead of overwriting the same location on every run, the upstream process creates a new partition with an automatically generated ID. Each run therefore gets its own location rather than replacing the object that Spark may already have listed.
Conceptually, the layout changed from something like:
Code:
.../filename.csv
to:
Code:
.../<generated-partition-id>/filename.csv
Once Spark starts reading a partition, that partition is not being replaced by the next run. The next run writes to a different partition.
After that change, the short reads stopped.
There are other ways to avoid the same race. A staging copy, a completion protocol, or a table format with committed snapshots can all provide a stable input. But those are alternatives, not what I did here. For this job, giving every run its own partition was enough.
The part that initially confused me was the S3 consistency guarantee. Every S3 request can be strongly consistent and this failure can still happen, because Spark does not perform the listing and all subsequent reads as one atomic operation.
I collected all the evidence I could from the affected runs. It narrowed the problem considerably, but it did not distinguish conclusively between the two remaining mechanisms.
What I can say from the result is simpler: once the runs stopped sharing a mutable input location, the problem stopped.