How does a file inbound channel adapter work, and how do you stop it from re-processing files it has already picked up?
answer
- FileReadingMessageSource + poller emits Message<File>
- Dedup lives in FileListFilters, not the source
- AcceptOnce = in-memory, forgets on restart
- Persistent AcceptOnce = MetadataStore-backed, survives restart/cluster
- LastModified/rename guards against half-written files
basics
~20 sA file inbound channel adapter uses FileReadingMessageSource, polled on a schedule, to emit each new file as a Message. To avoid re-reading the same file, add a FileListFilter such as AcceptOnceFileListFilter (in-memory) or a persistent one backed by a MetadataStore.
solid answer
~40 sThe file inbound channel adapter wraps a FileReadingMessageSource that scans a directory. A poller invokes it on a schedule; each qualifying file is emitted as a Message<File>. It does not remember files by itself — deduplication is the job of FileListFilters. AcceptOnceFileListFilter tracks seen files in memory, so it forgets on restart. For durability across restarts you use FileSystemPersistentAcceptOnceFileListFilter backed by a MetadataStore (e.g. PropertiesPersistingMetadataStore, or Redis/JDBC stores in a cluster). You typically also add a filter to ignore incomplete files still being written (e.g. filter by suffix, or a LastModifiedFileListFilter/ChainFileListFilter) and use a locker in multi-instance setups. Filters run before the file is emitted; without one you risk repeatedly processing the same file every poll.
code
java · 21 lines@Bean
public MetadataStore metadataStore() {
return new PropertiesPersistingMetadataStore(); // survives restarts on one host
}
@Bean
public IntegrationFlow csvIntakeFlow(MetadataStore metadataStore) {
ChainFileListFilter<File> filters = new ChainFileListFilter<>();
filters.addFilter(new SimplePatternFileListFilter("*.csv")); // only .csv
filters.addFilter(new LastModifiedFileListFilter()); // not still-being-written
filters.addFilter(new FileSystemPersistentAcceptOnceFileListFilter( // durable dedup
metadataStore, "csv-intake"));
return IntegrationFlow
.from(Files.inboundAdapter(new File("/data/in"))
.filter(filters)
.autoCreateDirectory(true),
e -> e.poller(Pollers.fixedDelay(1000).maxMessagesPerPoll(10)))
.handle("csvService", "process")
.get();
}go deeper
Knows the adapter reads files as messages and that some filter is needed to avoid re-reading.
Must name FileReadingMessageSource + poller and distinguish in-memory AcceptOnce from persistent MetadataStore-backed dedup.
Handles incomplete-file races, filter ordering in a chain, and clustered dedup with shared stores and lockers.
Designs end-to-end intake reliability: atomic-rename contracts with producers, at-least-once vs exactly-once semantics, store choice for HA, and scan strategy for very large directories.
**Anatomy.** A file inbound channel adapter (`<int-file:inbound-channel-adapter>` or `Files.inboundAdapter(dir)` in the DSL) is a **polled** inbound adapter backed by `FileReadingMessageSource`. Because it is polled, it needs a **poller** (`Pollers.fixedDelay(...)` / `@Poller`). On each poll cycle the source scans the directory, applies its filters, and, if any file qualifies, returns one `Message<File>` (payload is a `java.io.File`). Internally it uses a queue: a directory scan can enqueue many files, and subsequent polls drain the queue one message at a time (bounded by the poller's `maxMessagesPerPoll`). **Why duplicates happen.** The source re-scans the directory every poll. If a file is still present after processing (you didn't move/delete it), a naive source would emit it again and again. Deduplication is delegated to **`FileListFilter<File>`** implementations, not baked into the source. **The filters.** - `AcceptOnceFileListFilter` — remembers each file it has passed **in memory** (by name). Prevents re-emitting within a single run, but the memory is lost on **restart**, so files can be re-processed after a bounce. - `FileSystemPersistentAcceptOnceFileListFilter` — the durable version; it records seen files (name + last-modified) in a **`MetadataStore`**. With `PropertiesPersistingMetadataStore` the state survives restarts on one host; with a shared store (`RedisMetadataStore`, `JdbcMetadataStore`, `ZookeeperMetadataStore`) multiple app instances share dedup state. - `SimplePatternFileListFilter` / `RegexPatternFileListFilter` — accept only names matching a glob/regex (e.g. `*.csv`). - `LastModifiedFileListFilter` — accept only files older than N seconds, a simple guard against reading a file still being written. - `CompositeFileListFilter` / `ChainFileListFilter` — combine several filters. `ChainFileListFilter` is order-sensitive and stops passing a file down the chain once one filter rejects it, which is important so an *accept-once* filter doesn't mark a still-being-written file as seen. **Incomplete files.** A classic gotcha: a producer writes `data.csv` in place; the adapter reads it mid-write. Mitigations: producers write to a temp name then atomically rename; or add a `LastModifiedFileListFilter`; or watch for a companion marker file. Order the chain so the completeness check runs **before** the accept-once filter, otherwise the partial file gets recorded as seen and the completed version is skipped. **Concurrency / clustering.** With multiple instances polling the same directory, two can grab the same file. Use a `FileLocker` (e.g. `NioFileLocker`) via `.locker(...)`, and a **shared** `MetadataStore`-backed persistent filter, so exactly one instance processes each file. **Post-processing.** Filters prevent re-emission, but housekeeping (move/delete processed files) is usually done downstream with a `FileWritingMessageHandler` or a service activator, or by `preventDuplicates(false)` combined with delete-on-success. If you delete processed files, an in-memory accept-once filter is often enough; if you leave them in place, you need the persistent filter. **Alternatives.** `WatchServiceDirectoryScanner` uses OS file-system events instead of full re-scans for large directories. This changes scanning efficiency but the dedup story (filters + metadata store) is the same. **When to use what.** Short-lived, files-deleted-on-success, single instance → `AcceptOnceFileListFilter`. Files left in place and/or restarts matter → `FileSystemPersistentAcceptOnceFileListFilter` + `PropertiesPersistingMetadataStore`. Clustered → the persistent filter over a shared store plus a locker.
- Why does AcceptOnceFileListFilter re-process files after an application restart?It keeps the set of seen files only in memory (a Set on the heap). A restart wipes that memory, so on the next scan every still-present file looks new. Use FileSystemPersistentAcceptOnceFileListFilter over a MetadataStore to persist the state.
- How do you stop the adapter from reading a file that is still being written?Have producers write to a temp name and atomically rename when done, and/or add a LastModifiedFileListFilter (only accept files older than N seconds) or a marker-file check. Order it before the accept-once filter so the partial file isn't recorded as seen.
- How do you make dedup work across several app instances polling the same folder?Back the persistent filter with a shared MetadataStore (Redis/JDBC/Zookeeper) and add a FileLocker so only one instance claims each file.
saying these in an interview costs you the question
- Believing FileReadingMessageSource remembers processed files on its own
- Using AcceptOnceFileListFilter and expecting durability across restarts
- Ignoring the half-written-file race
- Putting the accept-once filter before the completeness check in the chain
- Assuming a single in-memory filter is safe across multiple instances