В статье мы разберем пример использования подхода — дробления данных (chunking). В мире разработки полно задач, связанных с обработкой больших объёмов данных. Будь то логи, дампы баз данных или файлы с медиа-контентом — эффективная работа с такими ресурсами становится критичной. В этой статье мы разберём практический пример многопоточной обработки файла в Java, где исходный файл разбивается на небольшие фрагменты (чанки) фиксированного размера. Такой подход лежит в основе многих систем: от торрент-клиентов до распределённых файловых хранилищ.
Постановка задачи
Предположим, у нас есть большой файл. Нам нужно разбить его на части по 16 килобайт каждая, причём каждая часть должна сохраняться в отдельный файл с суффиксом .moduleN, где N — порядковый номер чанка. Ключевые требования:
- Использовать многопоточность для ускорения обработки.
- Обеспечить конкурентное чтение из одного исходного файла без блокировок.
- Сохранять фрагменты параллельно.
- Корректно подсчитывать количество реально созданных модулей (даже при возможных ошибках).
В качестве демонстрационного кода возьмём класс Solution, который использует каналы (FileChannel) и пул потоков, а также счётчик AtomicLong для надёжного учёта.
Анализ кода
Структура и константы
|
1 |
private static final int CHUNK_SIZE = 16348; |
Размер чанка выбран равным 16 КБ — это кратно стандартному 4 КБ размеру блока во многих файловых системах, что позволяет эффективно использовать кэш и минимизировать накладные расходы на операции ввода-вывода.
Основной поток выполнения
- Считываем путь к файлу из консоли.
- Определяем размер файла и вычисляем количество чанков:
|
1 |
long chunkCount = fileSize / CHUNK_SIZE + (fileSize % CHUNK_SIZE > 0 ? 1 : 0); |
- Создаём пул потоков с фиксированным числом потоков (4). Это число выбрано исходя из количества ядер процессора (или логических потоков) — типичная практика для баланса между загрузкой CPU и накладными расходами на переключение контекста.
- Открываем один канал для чтения исходного файла до запуска задач. Этот канал будет использоваться всеми потоками.
- Для каждого чанка отправляем задачу в пул: чтение соответствующего сегмента файла и запись в новый файл.
- После отправки всех задач ожидаем их завершения с помощью
shutdown()иawaitTermination(). - Закрываем общий канал (автоматически благодаря
try-with-resources). - Выводим реальное количество успешно записанных модулей, сохранённое в
AtomicLong.
Ключевой элемент — FileChannel и transferTo()
Внутри каждой задачи используется FileChannel.open() для открытия исходного файла в режиме чтения и файла-чанка в режиме записи. Самый интересный момент — метод transferTo():
|
1 |
var savedLength = rd.transferTo(chunkIndex * CHUNK_SIZE, CHUNK_SIZE, wr); |
Этот метод позволяет скопировать данные непосредственно из одного канала в другой, используя возможности операционной системы (zero-copy). Вместо того чтобы читать данные в буфер Java-приложения, а затем записывать их, transferTo() делегирует работу ядру ОС, что существенно повышает производительность и снижает нагрузку на сборщик мусора.
Параметры:
- позиция в исходном файле — смещение, равное
номер_чанка * размер_чанка. - количество байт —
CHUNK_SIZE(для последнего чанка реальное количество может быть меньше). - целевой канал — файл для записи.
Метод возвращает реальное количество переданных байт, которое затем выводится в консоль.
Многопоточность и синхронизация
Надёжный подсчёт с AtomicLong
В коде используем счётчик:
|
1 |
AtomicLong savedChunkCount = new AtomicLong(0); |
При успешной записи каждого чанка выполняется:
|
1 |
savedChunkCount.getAndAdd(1); |
AtomicLong обеспечивает атомарное обновление без блокировок, что безопасно в многопоточной среде. В конце программы выводится именно это значение, а не расчётное chunkCount. Это даёт реальную информацию о том, сколько модулей действительно было создано — если в каком-то потоке произошла ошибка, счётчик не увеличится, и мы это увидим.
Пул потоков
Executors.newFixedThreadPool(4) создаёт ограниченное число потоков. Все задачи (чтение/запись чанков) будут выполняться параллельно, но не более четырёх одновременно. Это предотвращает чрезмерную конкуренцию за дисковый ввод-вывод, которая могла бы привести к деградации производительности из-за частых перемещений головки жёсткого диска (на HDD) или перегрузки шины (на SSD).
Потокобезопасность
В коде нет общих изменяемых данных между задачами, кроме самого файла, но чтение из одного файла несколькими потоками безопасно — ОС позволяет одновременное чтение. Каждый поток работает со своей позицией в файле, не пересекаясь с другими, поэтому синхронизация не требуется. Это главное преимущество: нет блокировок, все задачи могут выполняться независимо.
Ожидание завершения
|
1 2 |
pool.shutdown(); pool.awaitTermination(1, TimeUnit.HOURS); |
shutdown() запрещает приём новых задач, но уже отправленные продолжают выполняться. Затем основной поток блокируется до тех пор, пока все задачи не завершатся (или не истечёт таймаут в один час). Это гарантирует, что программа не завершится до окончания всех операций записи.
Исходный код
|
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 |
import java.io.*; import java.nio.channels.FileChannel; import java.nio.file.Files; import java.nio.file.Paths; import java.nio.file.StandardOpenOption; import java.util.Scanner; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicLong; public class Solution { // Фиксированный размер чанка (модуля) — 4 КБ private static final int CHUNK_SIZE = 16384; public static void main(String[] args) throws Exception { // Считываем путь к исходному файлу String filename; System.out.println("Введите путь к файлу:"); try (var sc = new Scanner(System.in)) { filename = sc.nextLine(); } var path = Paths.get(filename); long fileSize = Files.size(path); // счётчик созданных модулей long chunkCount = (long) fileSize / CHUNK_SIZE + (fileSize % CHUNK_SIZE > 0 ? 1 : 0); AtomicLong savedChunkCount = new AtomicLong(0); // Создаём пул из 4х потоков. ExecutorService pool = Executors.newFixedThreadPool(4); try (var rd = FileChannel.open(path, StandardOpenOption.READ)) { for (int i = 0; i < chunkCount; i++) { final int chunkIndex = i; pool.submit(() -> { String filenameNN = filename + ".module" + chunkIndex; try (var wr = FileChannel.open(Paths.get(filenameNN), StandardOpenOption.WRITE, StandardOpenOption.TRUNCATE_EXISTING, StandardOpenOption.CREATE)) { var savedLength = rd.transferTo(chunkIndex * CHUNK_SIZE, CHUNK_SIZE, wr); savedChunkCount.getAndAdd(1); System.out.println("Чанк №" + chunkIndex + " содержит " + savedLength + " байт!"); } catch (IOException e) { System.out.println("Ошибка записи чанка №" + chunkIndex); } }); } // Дожидаемся завершения операций. pool.shutdown(); pool.awaitTermination(1, TimeUnit.HOURS); } // Выводим общее количество созданных модулей System.out.println(savedChunkCount); } } |
Потенциальные проблемы и улучшения
Обработка ошибок
В коде есть обработка IOException внутри каждой задачи, но она лишь выводит сообщение. Для промышленного решения нужно предусмотреть повторные попытки, логирование и, возможно, прерывание всех задач при фатальной ошибке.
Порядок записи
Так как задачи выполняются параллельно, нет гарантии, что чанки будут записаны в порядке возрастания номеров. Но это не критично, если нам нужны только отдельные файлы. Однако если важна последовательность, можно использовать имена файлов с ведущими нулями для корректной сортировки.
Использование MappedByteBuffer
Для ещё более высокой производительности можно применить отображение файла в память (FileChannel.map()). Это позволит работать с данными как с массивом байтов в виртуальной памяти, но требует осторожности с размером файла и освобождением ресурсов.
Параллелизм vs. I/O
Разбиение файла на чанки — это I/O-интенсивная задача. Увеличение числа потоков сверх определённого предела не даёт выигрыша, так как узким местом становится диск. Оптимальное число потоков часто равно числу физических ядер или немного больше. В примере выбрано 4 — хороший компромисс для большинства случаев.
Вместо заключения
Данный пример демонстрирует способ параллельной обработки больших файлов в Java с использованием FileChannel и ExecutorService. Основные преимущества:
- Использование
transferTo()для эффективного копирования данных на уровне ОС. - Независимость потоков — отсутствие блокировок.
- Простота реализации и масштабирования.
Этот подход можно адаптировать для множества задач: параллельное шифрование, вычисление хешей, разбиение для загрузки по частям в облачное хранилище и многое другое. Понимание принципов работы с каналами и пулами потоков — важный шаг на пути к созданию высокопроизводительных приложений на Java.