Skip to content

CSVFile connector: sync() never deletes old snapshots, and fails or keeps the old token when the current file already has a snapshot #174

Description

@maximthomas

Summary

CSVFileConnector.sync() keeps every snapshot of the CSV file (<csv>.<mtime>) it has ever made, so syncFileRetentionCount has no effect. Both of its copies also mishandle a snapshot that already exists for the current file:

  • The copy at the start of a sync is made for a null token, or for a token whose snapshot is gone. If the snapshot already exists, it fails with ConnectorException: Unable to copy CSV file for sync operation.
  • The copy at the end of a sync that reported changes is skipped when the snapshot already exists. So is the token update, and the caller gets its old token back.

These have to be fixed together. Once old snapshots are deleted, consumers whose snapshot is gone will reach the start copy routinely.

Reported in #156 under "Related observations". Line references are to 9d53ed2 (current master). #170 rewrites the diff but not this code; see After #170.

Snapshots are never deleted

  • A sync that reports at least one change copies the whole CSV to <csv>.<mtime> under the read lock, unless that file already exists (L568-L580). A sync from a null token, or from a token whose snapshot is missing, makes its copy under the write lock before reading anything (L466-L482).
  • scrubSyncFiles (L959-L975) collects the snapshot names into a TreeMap and removes the oldest entries from that map until syncFileRetentionCount remain. It deletes no file. The setting defaults to 3 (CSVFileConfiguration.java#L68) and is described as "Number of sync history files to retain" (Messages.properties#L30).
  • So every modification of the CSV that a sync picks up leaves another full copy of the file on disk, and nothing ever removes it. With a 2 GB CSV, that is 2 GB more per modification.

An existing snapshot of the current file breaks sync

The copy at the start of sync() (L472-L482) calls Files.copy without checking whether <csv>.<mtime> exists. If it does, the copy throws FileAlreadyExistsException, and the caller gets ConnectorException: Unable to copy CSV file for sync operation. Two cases trigger it:

  • A second sync() with a null token on an unchanged file. The first one created <csv>.<mtime>.
  • A sync(T) whose snapshot <csv>.T is gone, while the file has not changed since the last sync that reported a change. That sync created <csv>.<mtime>.

The copy at the end of sync() does check whether <csv>.<mtime> exists (L572), but the new token is set inside that same check (L572-L579). Take two consumers of the same file. One consumer syncs from T and creates <csv>.<mtime>. When the other consumer then syncs from T, it gets the same deltas, but handleResult receives T again. It gets those deltas again on every sync until the CSV file changes.

The check and the copy are not atomic either. Both run under the shared read lock, so two concurrent syncs can both find no file. The second copy then throws the same exception, after that sync has already delivered its deltas.

A missing snapshot silently resets the baseline

If the start copy succeeds, sync(T) with <csv>.T gone uses the current file as the baseline and diffs it against itself. It passes the file's mtime to handleResult as the new token. Every delete since T is lost without an error. On master, creates and updates still arrive, as UPDATE deltas inside the full resend of unchanged rows that #156 fixes. After #170 nothing is sent at all.

This path is already reached today. When no snapshot exists, getLatestSyncToken returns the current time (L615), and no snapshot exists for that token either. Changes made between getLatestSyncToken and the first sync() are lost.

Once retention deletes snapshots, two more cases reach this path:

  • a consumer resuming from an old stored token;
  • one of several consumers of the same file, once the others have produced more than syncFileRetentionCount newer snapshots.

Snapshot names also match other files

scrubSyncFiles (L963) and getLatestSyncToken (L603) list the directory of the CSV file. They match each entry with Pattern.compile("(" + name + ")(\\.[0-9]{13})$") and Matcher.find(), where name is the CSV file name. The name is not quoted, and the pattern is not anchored at the start. For users.csv it therefore also matches old_users.csv.1300000000000 and usersXcsv.1300000000000. A name with regex metacharacters fails the other way: users(1).csv never matches its own snapshots, and users[1.csv makes sync() and getLatestSyncToken throw PatternSyntaxException. All four cases were checked with jshell.

Today the false matches only let getLatestSyncToken return another file's timestamp. A sync from that token takes the missing-snapshot path. If the file is unchanged since its last snapshot, the sync fails with the copy error; otherwise it silently resets the baseline. Once scrubSyncFiles deletes files, it would also delete another CSV file's snapshots in the same directory.

After #170

At the head of #170 (8a55cfa) this code is unchanged:

Possible fixes

  • Make scrubSyncFiles delete the files it drops, oldest first. It must never delete the snapshot the current sync is about to read.
    • Scrubbing runs under the write lock at the start of sync(). That lock is static and keyed by the absolute path of the CSV file (L148, L214). So no other sync() in the same JVM and connector bundle reads a snapshot of that file while it is deleted. Syncs in other JVMs are not excluded.
    • Log a failed delete instead of throwing: Windows cannot delete a file that is still open.
    • Scrubbing cannot move after the copy at the end of sync(), because that copy runs under the shared read lock. Up to syncFileRetentionCount + 1 snapshots therefore stay on disk until the next sync, or + 2 if the snapshot that sync read is older than the ones kept.
  • Write each snapshot to a temporary file and rename it into place, as CSVFile connector: a row deleted by an external writer during sync() is never reported as DELETE #172 also suggests, so that an existing <csv>.<mtime> is always complete. Then:
    • at the start of a sync, use an existing <csv>.<mtime> as the baseline instead of copying again;
    • at the end of a sync, treat an existing <csv>.<mtime> as done and still pass its mtime to handleResult.
  • In both scrubSyncFiles and getLatestSyncToken, match snapshot names exactly: Pattern.quote(name) + "\\.([0-9]{13})" with matches().
  • For a non-null token whose snapshot is gone, there are two options:
    • Fail with a clear ConnectorException, so that an operator resets the stored token. Only this option makes the lost changes visible. It also requires getLatestSyncToken to return only tokens that have a snapshot, for example by taking a snapshot when none exists. Today it returns the current time, and a token taken from it would fail on the first sync.
    • Keep the silent reset, but log a warning.

Activity

  1. self-assigned this
    on Oct 7, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

Labels

bugSomething isn't workingconnector:csvfileCSV file connectorresource-leakMemory, class-loader, thread or handle leaks

Type

No type

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions