Skip to content

CSVFile connector: sync() rescans the CSV for every row (O(N²)) and reports unchanged rows as UPDATE #156

Description

@maximthomas

Summary

CSVFileConnector.sync() becomes very slow on large CSV files. The file is not loaded into memory; it is streamed row by row. The problem is that for every row of one file, sync() reopens the other file and scans it from the top. As a result, one sync costs O(N²) row parses. A second bug reports every unchanged row as an UPDATE, so every sync also hands the whole file to the caller.

A sync with one changed row takes ~2 s at 500 rows, ~6 s at 1,000 rows and ~25–30 s at 2,000 rows, and returns N UPDATEs each time (see Reproducer). The target deployment has CSV files of up to ~2 GB (~8M rows), with the connector running at -Xmx16g. Extrapolating the measured quadratic cost, one sync at that size would take on the order of 10^8 s, i.e. years.

Line references are to 1abfe74. The timings come from a throwaway harness; everything else comes from reading the code.

1. Quadratic lookup in sync()

sync() compares the current CSV with the snapshot taken at the previous token in two passes:

  • Deletes (L494-L506): for every snapshot row, findObjectInFile(null, uid) opens the current CSV and scans it from the top until the uid is found, or to the end if the uid is absent.
  • Creates/updates (L532-L542): for every current row, findObjectInFile(syncOrigin, uid) opens the snapshot and scans it the same way.

Every deleted or created row costs a full scan of the other file, and every row present in both files costs a scan up to its match. That adds up to ~N² row parses and 2N file opens per sync. For 100,000 rows that is on the order of 10^10 parsed rows. Reading each file once would take 2×10^5 row parses.

Each scanned row is also expensive. findObjectInFile builds a ConnectorObject from every row before it compares the uid (L882-L883). When headerPassword is set, that includes a GuardedString (L855). Each GuardedString costs a SecureRandom IV, Cipher.getInstance("AES/GCM/NoPadding") with an AES-GCM encryption (AesGcmEncryptor.java#L58-L64) and a SHA-1 hash (GuardedString.java#L272).

The read lock is held for the whole comparison, from L477 to L577. Until the sync finishes, create, update, delete and another sync() on the same file wait for the write lock. The lock is a non-fair ReentrantReadWriteLock, and new readers block while a writer is first in the queue. So once one write is waiting, searches block as well.

2. Unchanged rows are reported as UPDATE

objectsDiffer (L932-L948) returns true as soon as it finds a differing attribute. When no attribute differs, it falls through to:

return left.getUid().getUidValue().equals(right.getUid().getUidValue())
        && left.getObjectClass().equals(right.getObjectClass())
        && left.getName().getNameValue().equals(right.getName().getNameValue());

sync() calls it only for a pair of objects matched by uid. Both objects get ObjectClass.ACCOUNT and use the uid as their name, so this expression is always true. Every unchanged row therefore produces an UPDATE delta.

Consequences:

  • Every sync delivers every row of the file, not just the changes. This holds even when the file has not changed since the previous token: the token stays the same (L559-L571), and each poll delivers N UPDATEs again.
  • A sync with a null token copies the file and compares the copy with the file itself (L463-L473). That costs N² work and produces N spurious UPDATEs instead of zero.

The existing test misses this. SyncOpTest.syncTest (L129-L134) removes the expected deltas from a map by uid and asserts that the map ends up empty. It never checks for extra deltas or for delta types. The test passes while the connector emits DELETE vilo, UPDATE miso, CREATE fanfi and a spurious UPDATE rado for the unchanged row. The expected map also uses CREATE_OR_UPDATE for miso and fanfi (L229, L244), while the connector emits UPDATE and CREATE. A test that checks types must fix those expectations too.

Reproducer

  1. Generate a CSV with N rows (uid,firstName,lastName,password). Configure headerUid=uid and headerPassword=password. Use a fresh file for each N, because a second null-token sync on an unchanged file fails (see "null token or missing snapshot" below).
  2. Run sync() with a null token and keep the token passed to handleResult. Correct result: 0 deltas. The connector snapshots the file at its current mtime and diffs the file against that snapshot.
  3. Change lastName in one row. Set the file's mtime explicitly, e.g. Files.setLastModifiedTime(csv, FileTime.fromMillis(old + 5000)). On a filesystem with coarse timestamps, an edit within the same tick keeps the old mtime.
  4. Time sync() from the saved token. Correct result: 1 UPDATE.

Results at 1abfe74 (JDK 26, macOS, APFS):

Rows Step 2: null-token sync Step 4: sync after one change
500 2.6 s, 500 UPDATEs 1.7 s, 500 UPDATEs
1,000 6.6 s, 1,000 UPDATEs 6.4 s, 1,000 UPDATEs
2,000 26.7 s, 2,000 UPDATEs 28.8 s, 2,000 UPDATEs

Timings vary by about 15% between runs. Step 2 at 500 rows includes JIT warm-up, since it is the first sync in the JVM. Doubling N makes each sync roughly 4× slower. Extrapolating the quadratic term gives most of a day for one sync at 100,000 rows.

Proposed fix

  1. Make objectsDiffer return false when no attribute differs. This alone fixes bug 2. Make SyncOpTest.syncTest assert the exact deltas and their types: vilo DELETE, miso UPDATE, fanfi CREATE, and nothing for rado.

  2. Replace the per-row lookup with a hash join over three streaming passes. The index exists only within a single sync() call; nothing is cached across calls. It is a single HashMap<String, byte[]> that maps each current uid to the fingerprint of the first snapshot row with that uid.

    1. Read the current CSV and put every uid into the map with a null value.
    2. Read the snapshot. If a row's uid is not in the map, emit a DELETE. Otherwise, if its value is still null, store the row's fingerprint. If a fingerprint is already stored (the uid repeats in the snapshot), keep the first one. This matches the current first-match lookup.
    3. Read the current CSV again. If a row's value is null, emit CREATE. If its fingerprint differs, emit UPDATE. If the fingerprints match, emit nothing.

    The fingerprint is a SHA-256 over the row's values, looked up by column name in getHeader() order. Each value is length-prefixed, so cell boundaries are unambiguous, and null is encoded as a separate marker:

    • compareHeaders accepts a snapshot whose columns are permuted (L968-L982); such a snapshot gets the same fingerprint.
    • A null cell is encoded differently from "". Super CSV returns null for both ,, and ,"",, and newConnectorObject then adds no attribute. A whitespace-only cell becomes "" after trimming.
    • Comparing passwords through the digest is equivalent to the current comparison, because GuardedString.equals compares SHA-1 hashes of the clear text (GuardedString.java#L276-L288).

    ConnectorObjects, and therefore GuardedStrings, are built only for the deltas that are emitted. For an UPDATE, the snapshot row's object is never built; as today, the delta carries only the current object. After this change sync() no longer calls findObjectInFile or objectsDiffer, and they have no other callers, so both can be removed.

    Every pass reads the header properly. That also removes a quirk of the deletes pass: findObjectInFile(null, ...) uses the cached header but never consumes the header line of the current file, so that line is parsed as a data row on every scan (L874-L881).

    What stays as it is today:

    • delta order: deletes in snapshot order, then creates and updates in file order;
    • handling of duplicate uids in either file;
    • a row with an empty uid cell still fails the sync. Today such a row anywhere in either file breaks every sync, because findObjectInFile builds an object from every row it scans.

    The only behavioral difference concerns a snapshot row whose uid equals the uid column name (e.g. uid). Today such a row matches the header line of the current file and is never reported as deleted.

    Cost: ~3N row parses instead of ~N². Memory is roughly 150 bytes per row for short uids (a map node, the uid String and a 32-byte digest). That is about 1.2 GB at 8M rows, which fits in -Xmx16g.

Related observations (not covered by the fix above)

  • Snapshot copies are never deleted. scrubSyncFiles (L950-L966) only removes entries from a local TreeMap, so syncFileRetentionCount (default 3, CSVFileConfiguration.java#L68) has no effect.

    • Every sync that processed a change copies the whole CSV under the read lock (L565) and keeps the copy forever. With a 2 GB CSV, that leaves another 2 GB file on disk after each such sync.
    • A sync from a null token makes its copy under the write lock (L457), which blocks searches too.
    • The hash join does not change this, because the snapshot copy is part of the design. I plan to fix retention in a separate PR.
  • null token or missing snapshot. The copy in L463-L473 does not check whether csv.<mtime> already exists. A second null-token sync on an unchanged file fails, and so does a sync from a token whose snapshot is gone while csv.<mtime> already exists. The error is ConnectorException: Unable to copy CSV file for sync operation, caused by FileAlreadyExistsException.

    When the copy does succeed, a missing snapshot silently resets the baseline to the current file and returns its mtime as the new token. Today the deletes since the old token are already lost. The creates and updates still arrive, but only as part of the full resend caused by bug 2. Once bug 2 is fixed, nothing at all is sent, and every change since the token is lost without an error.

    A single consumer always holds the token of the newest snapshot (L559-L571). So once retention deletes old snapshots, this path is reached by a consumer resuming from an old token, or by several consumers syncing the same file. It should be fixed together with retention.

  • Handler stop. A handler that returns false only breaks out of the current loop (L500-L502). The creates/updates pass still runs and offers the handler more deltas. As soon as any delta has been accepted, changesProcessed is true (L503, L540), the token advances (L559-L571), and the deltas after the stop point are never delivered. The token stays where it was only if the handler rejects every delta it is offered, which after a false is at most one per pass.

Activity

  1. self-assigned this
    on Oct 6, 2026
  2. added 2 commits that reference this issue 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 connectorperformancePerformance and scalability fixes

Type

No type

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions