1. Иногда хочется обрабатывать файлы параллельно
Когда вы только начинаете программировать, кажется логичным: «ну файл же один за другим — прочитал, обработал, пошёл дальше». И это правда работает… пока файлов мало, и они маленькие. Но в реальности часто бывает сценарий «у меня папка логов / фото / отчётов / исходников» и нужно сделать одно и то же для каждого файла: посчитать размер, строки, проверить формат, вытащить метаданные, скопировать, переименовать, построить отчёт.
И вот тут появляется естественное желание: раз файлов много, то почему бы не делать несколько штук одновременно? Коррутины как раз и придуманы, чтобы «запускать много задач» без того, чтобы вручную плодить потоки, от которых потом хочется спрятаться под стол. Kotlin прямо описывает корутины как лёгкий механизм конкурентности, который позволяет приостанавливать выполнение без блокировки системных ресурсов.
Но есть подвох: «параллельно» не значит «как можно больше сразу». Файловая система и диск — не бесконечные. Поэтому в этой лекции будет два слоя: как распараллелить и как не перестараться.
2. launch и async: похожи снаружи, разные по смыслу
На этом месте новички обычно начинают подозревать, что программисты нарочно придумывают похожие слова, чтобы вы страдали. И… да, иногда похоже на правду.
Оба — launch и async — это coroutine builder’ы, то есть функции, которые запускают корутину. Kotlin в документации прямо перечисляет их как базовые строители корутин. Но смысл у них разный:
- launch — «запусти и неси ответственность как за Job»: это корутина для действия.
- async — «запусти и дай мне результат потом»: это корутина для результата, она возвращает Deferred<T>.
Чтобы мозгу было проще, держите маленькую табличку:
| Что нужно | Что использовать | Что возвращает | Когда ошибка проявится |
|---|---|---|---|
| «Просто сделай» | |
|
обычно «внутри» (или через обработчики) |
| «Сделай и верни значение» | |
|
при await() / awaitAll() |
Мини‑пример: запускаем действие через launch
import kotlinx.coroutines.launch
import kotlinx.coroutines.runBlocking
fun main() = runBlocking {
val job = launch {
println("Я просто что-то делаю") // Я просто что-то делаю
}
job.join()
}
Мини‑пример: запускаем вычисление и ждём результат через async
import kotlinx.coroutines.async
import kotlinx.coroutines.runBlocking
fun main() = runBlocking {
val d = async { 2 + 3 }
println(d.await()) // 5
}
В этой лекции нас интересует именно async, потому что при обработке файлов нам почти всегда нужен результат: размер файла, статистика, текст, хэш, отчёт — что угодно.
3. Практический пример: утилита File Reporter
Чтобы примеры не были набором разрозненных заклинаний, представим, что мы пишем консольную утилиту File Reporter. Она получает список файлов и строит короткий отчёт: имя, размер, количество строк (для текстовых файлов).
Мы будем держать модель данных простой, без ООП‑фокусов «на вырост». Просто data class, потому что нам нужно удобно хранить результат.
import java.io.File
data class FileReport(
val file: File,
val sizeBytes: Long,
val lineCount: Int
)
Теперь сделаем suspend‑функцию, которая обрабатывает один файл. Важно: чтение файла — в Dispatchers.IO, а внутри используем .use {}.
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.withContext
import java.io.File
suspend fun buildReport(file: File): FileReport =
withContext(Dispatchers.IO) {
val size = file.length()
val lines = file.bufferedReader().use { it.lineSequence().count() }
FileReport(file = file, sizeBytes = size, lineCount = lines)
}
Да, lineSequence().count() внутри выглядит довольно «умно», но по смыслу там всё просто: читаем строки и считаем.
4. Шаблон: много файлов → много async → один awaitAll()
Самая приятная часть сегодняшней лекции: когда вы однажды увидели этот шаблон, вы начинаете замечать его в коде везде. Он похож на «конвейер»: мы создаём пачку задач и затем ждём, когда они все закончатся.
Технически это выглядит так:
- у нас есть List<File>
- мы делаем map { async { ... } } и получаем List<Deferred<FileReport>>
- вызываем awaitAll() и получаем List<FileReport>
awaitAll() удобен тем, что он «собирает результаты пачкой», и это сильно упрощает код по сравнению с ручным for и await().
Пример: параллельно построить отчёты по списку файлов
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.async
import kotlinx.coroutines.awaitAll
import kotlinx.coroutines.runBlocking
import java.io.File
fun main() = runBlocking {
val files = listOf(File("a.txt"), File("b.txt"), File("c.txt"))
val reports = files
.map { f -> async(Dispatchers.IO) { buildReport(f) } }
.awaitAll()
println("reports=${reports.size}") // reports=3
}
Обратите внимание на маленькую деталь: мы передали Dispatchers.IO прямо в async(...). Это означает, что корутина для каждого файла будет выполняться в I/O‑контексте (то есть блокирующее чтение файла допустимо). Это укладывается в общую модель «файлы в IO», которую мы приняли раньше.
И да: документация Kotlin подчёркивает, что launch() и async() — это базовые строители корутин, которые запускают конкурентные задачи внутри CoroutineScope.
5. Почему «тысяча async» — не всегда праздник
В какой-то момент студент делает так: «У меня 20 000 файлов. Я просто сделаю .map { async { ... } }.awaitAll()». И в этот момент компьютер начинает звучать как маленький аэропорт, а вы узнаёте, что такое «деградация производительности от слишком большого энтузиазма».
Проблема тут не в корутинах — корутины действительно лёгкие. Проблема в том, что:
- каждое чтение файла всё равно выполняет блокирующие вызовы ОС;
- файловая система и диск не ускорятся от того, что вы попросили «20 тысяч операций одновременно»;
- вы создаёте нагрузку на планировщик задач, очереди I/O, кеши и всё, что между вами и диском.
Поэтому мы делаем следующее «инженерное» движение: ограничиваем параллелизм.
Dispatchers.IO.limitedParallelism(n)
Идея простая: мы берём I/O‑диспетчер и говорим: «дорогой диспетчер, пожалуйста, не выполняй больше n задач одновременно». Остальные пусть подождут своей очереди.
import kotlinx.coroutines.Dispatchers
val io4 = Dispatchers.IO.limitedParallelism(4)
Затем запускаем все async на io4.
Пример: ограничение одновременности до 4 файлов
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.async
import kotlinx.coroutines.awaitAll
import kotlinx.coroutines.runBlocking
import java.io.File
fun main() = runBlocking {
val io4 = Dispatchers.IO.limitedParallelism(4)
val files = (1..10).map { File("file-$it.txt") }
val reports = files
.map { f -> async(io4) { buildReport(f) } }
.awaitAll()
println("done=${reports.size}") // done=10
}
Заметьте, что файлов 10, а параллелизм 4. Это не значит, что мы обработаем только 4 файла. Это значит, что одновременно «в работе» будет максимум 4 задачки, а остальные будут ждать, пока освободится слот.
Как выбрать n
Здесь нет «вечного правильного» числа, и это нормально. В учебном проекте обычно достаточно 2–8. Если поставить 64 «потому что красиво», вы, скорее всего, не получите ускорения, а иногда даже замедлите работу. В реальных проектах параллелизм подбирают экспериментально и/или конфигурируют.
6. Граница I/O и CPU: почему это влияет на скорость
Когда вы читаете файл, вы делаете I/O. Но часто после чтения начинается обработка: парсинг, поиск, подсчёты, нормализация текста, вычисление хэша. И это уже CPU‑работа.
Наивное решение: «всё сделаю внутри Dispatchers.IO, раз уж там файл». Работать будет, но архитектурно это нечётко: вы смешиваете разные типы нагрузки.
Нам полезно держать правило: внутри I/O‑границы делаем только то, что связано с файловой системой и потоками. Всё остальное — обычные вычисления, и им часто подходит Dispatchers.Default.
Пример: прочитать текст в IO, посчитать что-то в Default
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.withContext
import java.io.File
suspend fun countWords(file: File): Int {
val text = withContext(Dispatchers.IO) { file.readText() }
return withContext(Dispatchers.Default) {
text.split(Regex("\\s+")).count { it.isNotBlank() }
}
}
Этот пример специально короткий и не идеальный по производительности (регулярки и split могут быть дорогими), но он показывает архитектурную границу: чтение — отдельно, обработка — отдельно. И это делает программу предсказуемее.
7. Сборка отчётов: параллельно и безопасно
Теперь оформим «сердце» нашей лекции — функцию, которую можно будет спокойно вызывать из main. Мы сделаем её suspend, чтобы она работала в корутинном мире без runBlocking внутри (это важная дисциплина).
Мы также добавим параметр parallelism, чтобы можно было менять ограничение одновременности.
Базовый вариант: buildReportsParallel(...)
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.async
import kotlinx.coroutines.awaitAll
import kotlinx.coroutines.coroutineScope
import java.io.File
suspend fun buildReportsParallel(
files: List<File>,
parallelism: Int
): List<FileReport> = coroutineScope {
val io = Dispatchers.IO.limitedParallelism(parallelism)
files
.map { f -> async(io) { buildReport(f) } }
.awaitAll()
}
Тут есть важный момент: мы используем coroutineScope { ... }. Это значит, что все созданные async — «дети» текущей корутины, и функция не завершится, пока все дети не завершатся. Это и есть structured concurrency в действии: задачи связаны по жизненному циклу, а не «улетели в космос и забылись». Kotlin описывает контекст и иерархию задач как основу структурированной конкурентности.
Использование из main
import kotlinx.coroutines.runBlocking
import java.io.File
fun main() = runBlocking {
val files = listOf(File("a.txt"), File("b.txt"), File("c.txt"))
val reports = buildReportsParallel(files, parallelism = 2)
for (r in reports) {
println("${r.file.name}: ${r.lineCount} lines") // например: a.txt: 10 lines
}
}
Если один файл «плохой»: как не завалить всё сразу
Сейчас будет важная (и немного философская) часть: конкурентность — это не только про «быстрее». Это ещё и про «а если один из них упал?».
Поведение по умолчанию у такого кода довольно строгое: если внутри одного async случилось исключение, оно «всплывёт» при await()/awaitAll(). То есть ошибка станет видимой там, где вы ждёте результат. И это, честно говоря, хорошая модель: вы либо получили все отчёты, либо узнали, что не получилось.
В рамках этой лекции мы не будем строить сложную стратегию частичных результатов (это отдельная большая тема), но покажем простой приём: если вы хотите, чтобы обработка остальных файлов продолжалась, вы можете внутри async перехватывать ошибку и возвращать null.
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.async
import kotlinx.coroutines.awaitAll
import kotlinx.coroutines.coroutineScope
import java.io.File
suspend fun buildReportsOrSkip(
files: List<File>,
parallelism: Int
): List<FileReport> = coroutineScope {
val io = Dispatchers.IO.limitedParallelism(parallelism)
files
.map { f ->
async(io) {
try { buildReport(f) } catch (e: Exception) { null }
}
}
.awaitAll()
.filterNotNull()
}
Это не «идеальная обработка ошибок», но как учебный мостик — очень полезно: вы видите, как ошибки влияют на сбор результатов и почему awaitAll() — это точка, где результат становится «обязательным».
8. Типичные ошибки при параллельной обработке файлов
Ошибка №1: запускать файловые операции на Dispatchers.Default.
Это выглядит безобидно: «ну оно же работает». Но так вы начинаете блокировать потоки, которые предназначены для вычислений. В итоге другая корутина, которая честно считает что-то CPU‑тяжёлое, вдруг начинает ждать, потому что вы заняли её поток чтением файла. Лекарство простое и скучное, как техника безопасности: файловые операции всегда отправляем в Dispatchers.IO.
Ошибка №2: создавать async внутри map, но забывать про await/awaitAll.
Это классический «я всё запустил, но ничего не произошло». На самом деле произошло — вы создали пачку Deferred, но так и не дождались результатов. В лучшем случае программа завершится раньше времени, в худшем вы получите странные эффекты, потому что работа ещё идёт, а вы уже печатаете «готово». Лечится дисциплиной: если использовали async, значит где-то должен быть await() или awaitAll().
Ошибка №3: запускать слишком много задач одновременно без ограничения.
Когда файлов много, соблазн велик: «корутины же лёгкие». Да, но диск не резиновый. Если вы запускаете тысячи задач чтения одновременно, вы создаёте конкуренцию за ресурс, который и так один (ну ладно, два, если у вас два диска). Часто это замедляет выполнение и делает систему менее отзывчивой. Хорошая привычка — ставить limitedParallelism(n) и начинать с маленького n, например 4.
Ошибка №4: делать async(Dispatchers.IO) { withContext(Dispatchers.IO) { ... } }.
Это выглядит как «двойная защита от дурака», но на практике это просто лишняя вложенность и шум. Вы либо запускаете корутину на IO‑диспетчере, либо переключаетесь на IO внутри. В нашем стиле сегодня удобнее: async(ioLimited) { ... }, а внутри уже просто вызываем код работы с файлами (или функции, которые сами внутри переключаются — но тогда не надо второй раз).
Ошибка №5: смешивать чтение файла и тяжёлую обработку в одном IO‑контексте.
Такой код часто появляется сам собой: «я уже в IO, значит сделаю тут всё». Но IO‑контекст предназначен для блокирующих операций, а не для «долго считаем регулярками и парсим гигабайт текста». В результате вы забиваете I/O‑пул вычислениями и теряете предсказуемость. Гораздо чище: I/O отдельно, CPU отдельно — хотя бы двумя шагами.
ПЕРЕЙДИТЕ В ПОЛНУЮ ВЕРСИЮ