Папка как интеграция: файловый обмен в .NET без cron, FileSystemWatcher и самописных костылей
Спросите у любого интегратора, через что реально ходят деньги, и он назовёт файлы. Банк кладёт выписку в каталог на шаре. 1С выгружает заказы в XML по расписанию. Партнёр по EDI шлёт .edi на FTP. Биллинг ночью роняет CDR-файлы на пару гигабайт. Розничная сеть присылает остатки в CSV, потому что так было в 2009-м и никто не собирается это менять. Файловый обмен не «легаси, которое скоро уйдёт», это работающий транспорт, у которого есть ровно одно свойство: он выглядит проще, чем есть.
Проще он выглядит потому, что первая версия пишется за час. Directory.GetFiles, цикл, File.ReadAllBytes, обработали, удалили. Дальше начинается настоящая жизнь. Файл подхватили, пока партнёр его ещё дописывал, и в базу уехала половина заказов. Обработка упала, файл уже удалён. Два инстанса сервиса подняли один и тот же файл и создали дубли платежей. Кто-то положил в каталог .tmp от антивируса. FileSystemWatcher пропустил события, когда каталог оказался на SMB-шаре. Приложение выкатили в момент обработки, и файл на 800 МБ остался в состоянии «уже не там, ещё не тут». Каждая из этих проблем решается отдельным костылём, и через год у вас есть свой недо-фреймворк на 2000 строк, который никто не хочет трогать.
redb.Route.File закрывает этот список готовым транспортом. Каталог становится обычным источником сообщений в маршруте, ровно как очередь Kafka или HTTP-эндпоинт: опрос с фильтром, ожидание, пока файл дописан, идемпотентность, стратегия «что делать после обработки», атомарная запись через временный файл. Тот же самый DSL, те же EIP-паттерны после From, та же телеметрия. Дальше разберём по шагам, как это выглядит и почему выигрывает у самописного поллера.
Файл за минуту
services.AddRedbRoute(route =>
{
route.Services.AddRedbRouteFile();
route.AddRouteBuilder<MyRoutes>();
});
using redb.Route.File.Fluent;
// Читаем входящие CSV и складываем разобранное в очередь
From(FileDsl.Read("/data/incoming")
.Include("*.csv")
.MinAge(2000)
.SortBy("Modified")
.MoveTo("/data/processed"))
.Log("Взяли ${header.redbFile.Name}, ${header.redbFile.Length} байт")
.To("direct://parse");
// Пишем результат на диск атомарно
From("direct://export")
.To(FileDsl.Write("/data/outgoing")
.FileName("${header.orderId}.json")
.TempPrefix(".tmp-"));
Всё. Пакет не тянет внешних зависимостей, регистрация одной строкой, схема file дальше доступна и во флюентном виде, и строкой URI.
Эндпоинт это строка или билдер
Каждый эндпоинт читается двумя равнозначными способами: типобезопасным билдером в коде и обычным URI, который кладётся в appsettings.json и меняется без пересборки. Компилируются они в одно и то же, билдер буквально собирает строку.
// Флюентно
From(FileDsl.Read("/data/incoming").Include("*.csv,*.xml").Recursive().Delete())
// Строкой (идентично)
From("file:///data/incoming?include=*.csv,*.xml&recursive=true&delete=true")
На Windows путь пишется как file:///C:/data/incoming, ведущий слэш снимается автоматически. Это мелочь, но именно на ней обычно спотыкаются, когда конфиг переезжает между машинами.
Главная проблема файлового обмена: файл, который ещё пишут
Ни одна файловая система не даёт вам события «запись закончена». Партнёр открыл поток, льёт 300 МБ по медленному каналу, а ваш поллер уже увидел имя в каталоге. Дальше вопрос не «случится ли», а «как быстро». Коннектор даёт шесть независимых механизмов, и в проде обычно комбинируют два.
Минимальный возраст. Самое дешёвое: не трогать файл, пока с последней записи не прошло N миллисекунд.
FileDsl.Read("/data/incoming").MinAge(5000)
Ожидание стабилизации размера. Стратегия Changed смотрит на размер и время изменения через заданный интервал и отдаёт файл, только когда он перестал расти на заданное время. Это правильный ответ для больших выгрузок, которые пишутся минутами.
FileDsl.Read("/data/incoming")
.ReadLock("Changed")
.ReadLockCheckInterval(1000) // как часто смотрим
.ReadLockMinAge(5000) // сколько держится неизменным
.ReadLockTimeout(120000) // сколько ждём максимум
Маркерный файл. Стратегия MarkerFile создаёт рядом имя.redbLock через FileMode.CreateNew, то есть атомарно на уровне ОС. Второй инстанс, который увидит маркер, файл пропустит. Это рабочий ответ на «два пода читают одну шару»: гонку выигрывает ровно один, проигравший даже не открывает файл.
FileDsl.Read("/mnt/share/in").ReadLock("MarkerFile")
Эксклюзивный захват. Стратегия FileLock открывает файл эксклюзивно и держит хендл всё время обработки. Пока он держится, никто другой файл не откроет ни на чтение, ни на запись. Тонкость, ради которой это стоит упомянуть: хендл берётся с разрешением на удаление, поэтому пост-обработка своего же маршрута спокойно удаляет или переносит файл, не дожидаясь снятия лока, а тело читается через уже открытый хендл, а не вторым открытием.
FileDsl.Read("/data/incoming").ReadLock("FileLock")
Переименование. Стратегия Rename уводит файл под служебное имя и работает с ним там. Если переименование не удалось, значит файл держит кто-то другой, и маршрут его не трогает. Приятный побочный эффект: пока файл в обработке, его нет в каталоге под исходным именем, так что чужой поллер по маске его просто не увидит.
FileDsl.Read("/data/incoming").ReadLock("Rename")
Файл-сигнал от отправителя. Классика EDI: партнёр кладёт order.csv, дописывает, и только потом кладёт пустой order.csv.done. Пока сигнала нет, файла для маршрута не существует.
FileDsl.Read("/data/incoming").DoneFileName("${file:name}.done")
Поддерживаются подстановки ${file:name} и ${file:name.noext}, так что схема order.csv плюс order.done тоже собирается. После успешной обработки файл-сигнал удаляется сам.
Отдельно: консюмер по умолчанию игнорирует служебные имена, начинающиеся с .redb_, а также всё, что начинается с настроенного tempPrefix. То есть продюсер, пишущий в тот же каталог через временный файл, не устроит сам себе бесконечный цикл. Это та мелочь, которую в самописном поллере обнаруживают на третий день.
Что делать с файлом после обработки
Стратегия ровно одна на эндпоинт, и попытка задать две сразу падает на валидации при старте, а не молча в рантайме.
| Стратегия | Что происходит | Когда брать |
|---|---|---|
Noop() |
Файл остаётся на месте | Каталог только читаем, права на запись может не быть |
Delete() |
Файл удаляется после успеха | Входящая «мусорка», где важен только факт приёма |
MoveTo(dir) |
Файл переезжает в архив | Продакшн по умолчанию: всегда есть что предъявить |
PreMove(dir) |
Файл переезжает ДО обработки | Несколько потребителей на одном каталоге |
PreMove заслуживает пояснения, потому что решает задачу, которую обычно решают неправильно. Файл сначала переносится в рабочий подкаталог, и только потом читается. Перенос внутри одной файловой системы атомарен, поэтому второй инстанс просто не найдёт файл и пойдёт дальше. Плюс вы в любой момент видите по каталогу, что именно сейчас в работе.
FileDsl.Read("/data/incoming")
.PreMove("/data/processing")
.MoveTo("/data/archive")
Отдельно про неудачу. Если обработка упала и исключение не помечено обработанным, пост-обработка не выполняется: файл остаётся ровно там, где был, и попадёт в следующий опрос. Это осознанное поведение, потому что «удалить файл, который не смогли обработать» является худшим из возможных вариантов. Обратная сторона в том, что заведомо битый файл будет пробоваться снова и снова, поэтому на маршруте с файлами OnException не опция, а часть конструкции:
OnException(typeof(FormatException))
.Handled(true)
.To(FileDsl.Write("/data/quarantine"));
Обработанное исключение закрывает обмен успешно, файл уезжает в карантин копией, а оригинал обрабатывается штатной стратегией.
Сюда же относится случай, о котором обычно не думают заранее: файл нашёлся в каталоге, но прочитать его не удалось. Кто-то держит его эксклюзивно, права не те, том отвалился. Это тоже отказ по файлу, а не пустое сообщение: маршрут не получит ноль байт под видом успешной обработки, файл останется на месте, и его возьмут на следующем опросе. Звучит очевидно ровно до того момента, когда самописный поллер отдаёт вам пустой byte[], а импорт честно записывает в базу ноль заказов.
Идемпотентность
Включается флагом. Ключ по умолчанию собирается из полного пути, времени изменения и размера, так что перезаписанный под тем же именем файл считается новым, а тот же самый файл повторно не поедет.
FileDsl.Read("/data/incoming").Noop().Idempotent()
Связка Noop() + Idempotent() даёт режим «читаем каталог, ничего в нём не меняем, каждый файл ровно один раз». Это единственный способ работать с каталогом, куда у вас нет прав на запись, и он же самый частый при интеграции с чужой шарой.
Реестр обработанных ключей живёт в памяти процесса, и это надо держать в голове: после рестарта он пуст. В связке с Delete или MoveTo это не имеет значения, потому что обработанного файла в каталоге уже нет. В связке с Noop рестарт означает повторный проход по каталогу, и защиту от дублей в этом сценарии нужно ставить ниже по маршруту.
Отбор: что берём и в каком порядке
Порядок обработки в файловом обмене почти всегда значим. Ночная пачка из 4000 файлов, которую разбирают в случайном порядке, даёт остатки склада, разъезжающиеся с реальностью.
FileDsl.Read("/data/incoming")
.Include("*.csv,*.xml") // маски через запятую, * и ?
.Exclude("*.tmp,~*")
.SortBy("Modified") // Name, NameDesc, Modified, ModifiedDesc, Size, SizeDesc
.MaxMessagesPerPoll(200) // порция за один проход
.Recursive()
.Delay(1000) // интервал опроса, по умолчанию 500 мс
MaxMessagesPerPoll в паре с сортировкой это ваш регулятор нагрузки. Пришло 40 000 файлов после суточного простоя партнёра, а вы разбираете их предсказуемыми порциями от старых к новым, вместо того чтобы построить в памяти список на 40 000 элементов и уронить процесс.
Маски сравниваются без учёта регистра, что на Linux спасает от классического «прислали ORDERS.CSV, а маска *.csv».
Запись: temp плюс rename, а не «как получится»
Продюсер по умолчанию решает зеркальную задачу. Если вы пишете файл прямо в каталог, который читает партнёр, он гарантированно однажды прочитает половину.
From("direct://export")
.To(FileDsl.Write("/data/outgoing")
.FileName("orders-${header.batchId}.json")
.TempPrefix(".tmp-")
.AutoCreate());
Тело пишется в .tmp-orders-42.json, и только после полной записи файл переименовывается в целевое имя. Переименование внутри файловой системы атомарно, поэтому партнёр видит либо ничего, либо готовый файл целиком. Промежуточного состояния не существует. При исключении временный файл убирается за собой.
Имя целевого файла берётся, в порядке приоритета, из опции FileName (с полноценными выражениями по заголовкам и телу), из входящего заголовка redbFile.Name, а если нет ни того ни другого, генерируется как redb-{guid}. Последнее удобнее, чем кажется: маршрут From(kafka).To(file) начинает работать без единой настройки имени.
Что делать, если целевой файл уже есть, задаётся отдельно:
FileExist |
Поведение |
|---|---|
Override |
Перезаписать (по умолчанию) |
Append |
Дописать в конец |
Fail |
Бросить исключение |
Ignore |
Тихо пропустить запись |
Move |
Отложить существующий в .bak, затем писать |
TryRename |
Переименовать существующий, добавив метку времени |
Ignore и Fail выглядят экзотикой ровно до первой выгрузки, где повторная запись того же имени означает повторную отгрузку товара.
Имя файла это недоверенный ввод
Обратите внимание на порядок из предыдущего абзаца: если FileName не задан, имя приезжает из заголовка redbFile.Name. А заголовок туда положил кто-то: входящий файл партнёра, HTTP-загрузка, поле сообщения из очереди. То есть в типичном маршруте From(http).To(file) имя создаваемого файла выбирает тот, кто шлёт запрос.
Поэтому продюсер по умолчанию не выпускает запись за пределы своего каталога. Имя вида ../../etc/cron.d/backdoor или абсолютный путь дают исключение, а не файл в неожиданном месте. Отдельно стоит знать про абсолютный путь: обычный Path.Combine в .NET на нём молча выбрасывает базовый каталог и возвращает то, что ему передали, так что «мы же склеиваем с базой» защитой не является.
// по умолчанию: пишем только внутрь /data/outgoing
FileDsl.Write("/data/outgoing")
// осознанно разрешаем выход наружу
FileDsl.Write("/data/outgoing").JailStartingDirectory(false)
Опция называется одинаково у локальной ФС, FTP и SFTP, так что правило одно на все три транспорта.
Журнал в файл через Append
Отдельный полезный режим: append с разделителем. Приёмный лог, аудит-трейл, накопительный CSV за день собираются одним шагом маршрута.
From("direct://audit")
.To(FileDsl.Write("/var/log/app")
.FileName("audit-${dateFormat(now(), 'yyyyMMdd')}.log")
.FileExist("Append")
.AppendChars("\n"));
Имя с датой в выражении означает, что ротация по суткам получается сама собой, без отдельного планировщика.
Заголовки: имя файла едет через весь маршрут
Консюмер кладёт в сообщение полный набор метаданных с префиксом redbFile., доступный дальше в любом выражении, предикате или процессоре.
| Заголовок | Что внутри |
|---|---|
redbFile.Name |
order-42.csv |
redbFile.NameOnly |
order-42 |
redbFile.Extension |
.csv |
redbFile.AbsolutePath |
Полный путь на момент чтения |
redbFile.RelativePath |
Путь относительно каталога опроса |
redbFile.Parent |
Родительский каталог |
redbFile.Length |
Размер в байтах |
redbFile.LastModified |
DateTimeOffset |
Продюсер после записи проставляет redbFile.NameProduced с фактическим путём, куда файл лёг. Плюс к этому по расширению выставляется ContentType для .json, .xml, .csv, .txt, .log, .html, так что следующий шаг маршрута сразу знает, что ему прислали.
Практический эффект виден на маршрутизации по имени, а её в файловом обмене больше, чем хотелось бы:
From(FileDsl.Read("/data/in").Recursive().MoveTo("/data/archive"))
.Choice(c => c
.When(Header("redbFile.Extension").isEqualTo(".xml").Matches,
w => w.To("direct://xml"))
.When(Header("redbFile.Name").startsWith("INV_").Matches,
w => w.To("direct://invoices"))
.Otherwise(o => o.To("direct://unknown")));
Большие файлы: не читать целиком
По умолчанию тело это byte[], что удобно и совершенно неприемлемо для CDR-файла на два гигабайта. Режим StreamBody отдаёт открытый поток вместо массива, и поток закрывается вместе с обменом.
From(FileDsl.Read("/data/cdr").Include("*.dat").StreamBody().MoveTo("/data/done"))
.Split(ex => ReadLines((Stream)ex.In.Body!))
.To("direct://cdr-record");
Комбинация «стриминговое тело плюс сплиттер» это то, ради чего сплиттер в ESB и существует: файл разбирается построчно, память держит одну запись, а не весь файл.
Файловая система как камера хранения
Есть смежная задача, и она возникает ровно в файловых интеграциях. Тело большое, а промежуточные шаги маршрута его вообще не смотрят. Дальше по цепочке брокер, HTTP-вызов, ещё один сервис, и каждый честно тащит через себя двести мегабайт, которые ему не нужны.
Паттерн Claim Check (Хоуп и Вульф) решает это как гардероб: тело сдаётся в хранилище, а по маршруту едет номерок. Разворачивается тело там, где действительно нужно.
Хранилище подключаемое, и для файлового мира естественный выбор это сама файловая система: FileClaimCheckRepository кладёт каждую сдачу отдельным файлом рядом с метаданными и TTL. Для больших тел это дешевле памяти, а на общей шаре переживает и перезапуск, и переезд обработки на соседний инстанс.
private readonly IClaimCheckRepository _claims =
new FileClaimCheckRepository("/data/claims", TimeSpan.FromHours(6));
From(FileDsl.Read("/data/incoming").Include("*.zip").MoveTo("/data/archive"))
.ClaimCheck(_claims, ClaimCheckOperation.Set, "${header.redbFile.NameOnly}")
.To("direct://notify") // дальше едет номерок, не архив
.ClaimCheck(_claims, ClaimCheckOperation.GetAndRemove, "${header.redbFile.NameOnly}")
.Process(UnpackAndImport); // тело вернулось, номерок погашен
Тип исходного тела запоминается в заголовках, поэтому возвращается не голый массив байт, а то, что было. Хранилище можно зарегистрировать под именем в контексте и ссылаться строкой, а можно вообще не указывать: тогда шаги берут общее хранилище контекста.
Есть и стековый режим Push / Pop без ключа: убрали тело перед обогащением, вернули после, вложенность считается сама. Удобно, когда посреди маршрута нужен вызов, которому ваши двести мегабайт только мешают.
Локальная папка, FTP и SFTP это один и тот же маршрут
Здесь основная архитектурная ставка коннектора. Опрос, фильтры, сортировка, идемпотентность, doneFileName, пост-обработка, атомарная запись и стратегии существующего файла реализованы один раз в общей базе redb.Route.GenericFile. Локальная файловая система, FTP и SFTP это три реализации файловых операций поверх неё.
Практический смысл прямой. Когда партнёр говорит «мы больше не монтируем шару, забирайте с нашего SFTP», вы меняете источник, а не переписываете логику приёма:
// было
From(FileDsl.Read("/mnt/partner/in").Include("*.edi").MoveTo("archive"))
// стало
From(SftpDsl.Directory("/upload/in")
.Host("sftp.partner.com").Username("edi").Password("{{sftp-pass}}")
.Include("*.edi").MoveTo("archive"))
Имена опций, семантика пост-обработки и поведение при ошибке совпадают, потому что это буквально один и тот же код. Мигрировать между транспортами получается за минуту, а не за спринт.
Остановка без потери файла
Консюмер построен на общем для всех транспортов redb.Route механизме graceful shutdown: на Stop сначала прекращается опрос каталога, затем движок дожидается завершения уже начатых обменов и только после этого гасит маршрут. Файл, который в момент выката находился в обработке, не бросается на середине и не остаётся в подвешенном состоянии между PreMove и архивом.
Для файлового обмена это заметно сильнее, чем для брокеров: у брокера непонятое сообщение вернётся в очередь, а у файла второго шанса на автомате нет.
Границы
Честный список того, что коннектор не делает, чтобы вы не выясняли это в проде.
| Граница | Как есть |
|---|---|
| Опрос, а не подписка на события ФС | Никакого FileSystemWatcher. Опрос предсказуем, переживает сетевые шары и не теряет события при переполнении буфера. Цена: задержка до delay, по умолчанию 500 мс |
| Реестр идемпотентности в памяти процесса | Переживает рестарт только вместе с Delete или MoveTo, которые убирают файл из каталога. Он же не общий между инстансами: разводите их preMove или readLock, а не надеждой на общий реестр |
| Ключ идемпотентности собирается из метаданных файла | Путь, время изменения, размер, либо ваш шаблон с ${file:name}. Выражения по заголовкам и телу тут недоступны: ключ нужен до того, как файл прочитан, иначе проверка теряет смысл |
| Обработка внутри опроса последовательная | Один консюмер разбирает порцию файлов по очереди. Параллелизм набирается ниже по маршруту или несколькими маршрутами на разные маски |
Стратегии readLock только для локальной ФС |
У FTP и SFTP на их месте doneFileName, minAge и preMove |
| Каталог опроса не создаётся сам | Если каталога нет, опрос молча пропускает цикл. Автосоздание есть у продюсера через AutoCreate |
Куда это встраивается
Ценность файлового коннектора не в чтении каталога, это Directory.GetFiles. Ценность в том, что после From доступен весь остальной redb.Route: сплиттер и агрегатор, маршрутизация по содержимому, дедупликация, ретраи, circuit breaker, транзакции, распределённая трассировка. Типовой ночной приём выглядит так:
From(FileDsl.Read("/data/incoming")
.Include("orders_*.csv")
.DoneFileName("${file:name}.done")
.SortBy("Modified")
.MaxMessagesPerPoll(500)
.PreMove("/data/processing")
.MoveTo("/data/archive"))
.Log("Пачка ${header.redbFile.Name}")
.Split(ex => ParseCsv(ex.In.Body))
.To("sql:INSERT INTO orders(...) VALUES(...)?dataSource=#pg")
.End()
.To("kafka://orders-imported");
Двадцать строк вместо самописного сервиса, и каждая строка описывает решение, а не механику его исполнения. Продюсер, помимо прочего, открывает span OpenTelemetry на запись, так что «куда делся файл» перестаёт быть вопросом к логам.
Файловый обмен не станет модным. Но он останется, и разница между «у нас есть файловая интеграция» и «у нас есть надёжная файловая интеграция» измеряется ровно теми деталями, которые перечислены выше: возраст файла, маркер, атомарное переименование, порядок и порция. Приятно, когда их уже написали.
Пакет: redb.Route.File на NuGet; исходники и полный справочник по опциям в README коннектора. Файлы это ещё один транспорт в семействе redb.Route, рядом с Kafka, RabbitMQ, SFTP, AS2 и прочими: тот же From → … → To, те же EIP, та же наблюдаемость. Разница лишь в том, что на входе каталог, который партнёр наполняет тогда, когда ему удобно.
Если было полезно, ⭐ на GitHub поможет другим это найти.
Другие мои статьи — redb.ru/articles, ещё — на Хабре.