Ограниченные ресурсы Android-устройств часто воспринимаются как непреодолимое препятствие для задач «big data». Однако в реальной практике многие полезные кейсы (агрегации, предобработка логов, ETL небольших/средних выборок, аналитику с умеренным объёмом данных) можно выполнять на смартфоне или планшете, если корректно организовать запуск вычислительного движка и тщательно настроить потребление памяти.
В этом материале, как ведущий эксперт РыбинскЛАБ, покажу подход к обработке данных с помощью Apache Spark в режиме локального кластера прямо в Termux. Акцент — на оптимизации памяти, чтобы Spark и JVM оставались в рамках доступной RAM, а вычисления были стабильными.
Что именно будем запускать и почему это работает
Apache Spark рассчитан на распределённые вычисления, но умеет работать и локально: при запуске драйвера и исполнителей (executors) в пределах одного устройства. Это позволяет использовать знакомую модель вычислений (DataFrame/Dataset, SQL) и при этом избегать сетевой инфраструктуры.
Ключевая идея: минимизировать количество одновременно активных задач, контролировать размер кэшей/пулов и ограничивать память JVM так, чтобы система не «убивала» процесс.
Требования и подготовка Termux
Перед стартом проверьте:
- Свободное место (Spark и зависимости занимают заметный объём).
- Достаточную RAM (ориентируйтесь на 4–8 ГБ для комфортной работы; на меньших объёмах потребуется более жёсткая настройка памяти и небольшие датасеты).
- Устройство должно поддерживать нужную версию Java/JDK, доступную в Termux.
Подготовим окружение Termux. Пример (версии могут отличаться — проверяйте актуальные пакеты):
pkg update
pkg upgrade -y
pkg install -y openjdk-17 wget tar xz-utils procps coreutils
Проверьте Java:
java -version
Установка Apache Spark
Есть несколько способов установки Spark в Termux. Для практичности используем скачивание бинарного дистрибутива и распаковку.
1) Скачайте Spark (пример для 3.x; подбирайте версию под свою задачу):
SPARK_VERSION=3.5.1
wget -O spark.tgz https://archive.apache.org/dist/spark/spark-${SPARK_VERSION}/spark-${SPARK_VERSION}-bin-hadoop3.tgz
mkdir -p ~/spark
tar -xzf spark.tgz -C ~/spark --strip-components=1
rm -f spark.tgz
2) Проверьте структуру:
ls -la ~/spark
ls -la ~/spark/bin
3) Установите переменные окружения (можно добавить в ~/.bashrc или профиль Termux):
export SPARK_HOME=~/spark
export PATH=$SPARK_HOME/bin:$PATH
Концепция локального кластера в Spark
В локальном режиме Spark использует один JVM-процесс-драйвер (и созданные внутри того же устройства исполнители), но управляет задачами так, будто это кластер. Настройка строится вокруг параметров:
- Master: задаём
local[*]илиlocal[N]. - Количество исполнителей/параллелизм: ограничиваем, чтобы не раздувать потребление памяти.
- Память JVM: задаём upper bound через флаги
-Xmxи связанные параметры. - Очереди/буферы: контролируем размер shuffle и кэшей.
Оптимизация памяти: базовая стратегия
Чтобы Spark работал на ограниченном устройстве, обычно недостаточно лишь указать --executor-memory. Нужно учитывать:
- Память на уровне Spark и память JVM (heap/metaspace).
- Shuffle может резко увеличивать потребление памяти/диска.
- Кэширование
.cache()/.persist()ускоряет повторные вычисления, но может привести к OutOfMemory, если датасет не помещается.
Рекомендуемый подход:
- Начните с небольшого числа потоков:
local[2]илиlocal[3]. - Задайте явные лимиты heap для драйвера и исполнителей.
- Отключите агрессивное кэширование по умолчанию.
- Сделайте shuffle максимально «лёгким» за счёт сжатия и умеренных параметров количества партиций.
Запуск Spark-shell и настройка JVM/heap
Создадим рабочую сессию и зададим параметры памяти. Ниже пример для смартфона с ограниченными ресурсами. Корректируйте значения под своё устройство.
SPARK_HOME=~/spark
$SPARK_HOME/bin/spark-shell \
--master local[2] \
--conf spark.driver.memory=1g \
--conf spark.executor.memory=512m \
--conf spark.executor.cores=1 \
--conf spark.sql.shuffle.partitions=16 \
--conf spark.default.parallelism=2 \
--conf spark.serializer=org.apache.spark.serializer.KryoSerializer \
--conf spark.kryoserializer.buffer.max=64m \
--conf spark.shuffle.compress=true \
--conf spark.rdd.compress=true \
--conf spark.sql.inMemoryColumnarStorage.compressed=true \
--conf spark.sql.autoBroadcastJoinThreshold=10485760
Примечание: параметры spark.driver.memory и spark.executor.memory — это только часть картины. JVM heap ограничивается также флагами запуска. В локальном режиме и в зависимости от способа запуска Spark может отличаться распределение памяти, но в целом безопаснее держать разумные значения и не превышать доступную RAM.
Если вам нужно дополнительно настроить JVM-флаги (актуально при проблемах с heap), можно задать через конфигурацию Java-параметров. Практический вариант — формировать SPARK_SUBMIT_OPTS перед запуском.
export SPARK_SUBMIT_OPTS="-XX:+UseG1GC -XX:MaxGCPauseMillis=200 -XX:InitiatingHeapOccupancyPercent=35"
$SPARK_HOME/bin/spark-shell --master local[2] \
--conf spark.driver.memory=1g \
--conf spark.executor.memory=512m
Настройка путей, временных файлов и диска
В Spark активно используются временные директории для shuffle и промежуточных артефактов. На Android важно, чтобы хранилище было доступно и достаточно быстрое.
Рекомендуется явно задать рабочие директории (особенно если у вас много временных данных):
export TMPDIR=$HOME/tmp
mkdir -p $TMPDIR
$SPARK_HOME/bin/spark-shell \
--master local[2] \
--conf spark.local.dir=$TMPDIR \
--conf spark.driver.memory=1g \
--conf spark.executor.memory=512m \
--conf spark.sql.shuffle.partitions=16
Пример пайплайна: чтение данных, агрегации и фильтрация
Давайте рассмотрим типовую задачу: обработка JSON/CSV логов с последующей агрегацией. Важно избегать «широких» операций с большим количеством shuffle без необходимости.
Пример скрипта на Scala использовать в Termux можно через spark-shell, но проще — как базовый template. Ниже — пример для чтения CSV и агрегации (структуру полей подстройте под ваши данные).
// В spark-shell
import org.apache.spark.sql.functions._
val df = spark.read
.option("header", "true")
.option("inferSchema", "true")
.csv("file:///sdcard/your_path/data.csv")
val cleaned = df
.filter(col("status") === "ok")
.select(col("user_id"), col("event_time"), col("duration_ms"))
// Чуть умереннее с партициями, чтобы не раздувать shuffle
val result = cleaned
.groupBy("user_id")
.agg(
avg("duration_ms").as("avg_duration"),
count(lit(1)).as("events_cnt")
)
result.show(50, truncate=false)
// Осторожно с cache: используйте только если действительно нужно повторное обращение
// cleaned.cache().count()
Как уменьшить потребление памяти в Spark на Android
Ниже — практики, которые обычно дают максимальный эффект именно на «малых» устройствах:
- Ограничьте параллелизм:
--master local[2],spark.default.parallelismиspark.sql.shuffle.partitionsподберите под CPU и RAM. - Не кэшируйте лишнее:
.cache()и.persist()повышают риск OOM. Если кэш нужен — выбирайте разумные объёмы. - Снижайте объём shuffle: фильтруйте и проектируйте (операция select) как можно раньше, до широких преобразований.
- Используйте сжатие:
spark.shuffle.compress=true,spark.rdd.compress=trueи сжатое хранение колонок. - Следите за схемой данных: лишние строки/колонки увеличивают память и стоимость сериализации.
- Планируйте размер партиций: слишком много партиций создаёт overhead, слишком мало — приводит к большим задачам и росту памяти на executor.
Управление памятью во время работы: диагностика и типовые ошибки
На практике чаще всего встречаются:
- OutOfMemoryError в драйвере: обычно из-за слишком больших сборок (
collect,toPandas), чрезмерного кэша или слишком широкого shuffle. - GC overhead limit exceeded: heap слишком мал для текущей нагрузки или слишком активный аллокационный профиль.
- Слишком медленный shuffle: партиции/диск/сжатие настроены не оптимально.
Что делать при проблемах:
- Уменьшить
spark.sql.shuffle.partitions(в пределах разумного) и параллелизм. - Ограничить выборку до проверочного поднабора (например, по дате) и только потом запускать «полный» проход.
- Убедиться, что вы не собираете весь датасет в память через
collectили отображение огромных результатов.
Использование локальной сети (если нужна координация)
Иногда требуется, например, доступ к данным на другом устройстве или запуск взаимодействующих процессов в пределах вашей инфраструктуры. В таком случае корректный подход — настроить локальную сеть через VPN-туннель для локального взаимодействия, а не для обхода блокировок. Далее Spark при локальном master всё равно остаётся главным образом «локальным», но доступ к файлам/службам может быть проще организовать.
Если у вас есть конкретная топология (2 устройства, общий каталог, отдельный файловый сервер), опишите — помогу подобрать безопасную схему.
Работа с данными на диске: что важно для Termux
Termux работает с файловой системой Android, но важно учитывать права доступа и путь к данным. Для больших наборов лучше размещать файлы так, чтобы они были доступны процессу Spark и находились на носителе с приемлемой скоростью.
- Проверьте доступность пути (в примере использовался
file:///sdcard/...). - Убедитесь, что ваши CSV/JSON не приводят к гигантским схемам (особенно при
inferSchema=true). - Для повторных экспериментов используйте паркет (
parquet) как более эффективный формат, но стартуйте с ваших исходных данных.
Готовый набор конфигураций под «малую» RAM
Ниже — компактный набор конфигураций для старта. Он подходит для пробных задач и средних объёмов данных. В дальнейшем параметры стоит подстроить.
$SPARK_HOME/bin/spark-shell \
--master local[2] \
--conf spark.driver.memory=1g \
--conf spark.executor.memory=512m \
--conf spark.executor.cores=1 \
--conf spark.sql.shuffle.partitions=16 \
--conf spark.default.parallelism=2 \
--conf spark.serializer=org.apache.spark.serializer.KryoSerializer \
--conf spark.kryoserializer.buffer.max=64m \
--conf spark.shuffle.compress=true \
--conf spark.rdd.compress=true \
--conf spark.sql.inMemoryColumnarStorage.compressed=true \
--conf spark.sql.autoBroadcastJoinThreshold=10485760 \
--conf spark.local.dir=$HOME/tmp
Заключение
Запуск Apache Spark в режиме локального кластера в Termux на ограниченном устройстве — реалистичный и практичный сценарий, если вы подходите к настройке памяти и вычислительных параметров осознанно. Основные рычаги — ограничение параллелизма, контроль heap через spark.driver.memory/spark.executor.memory, настройка shuffle и осторожное использование кэшей. Такой подход позволяет выполнять полезные аналитические и ETL-задачи без полноценной серверной инфраструктуры.
Если хотите подобрать оптимальные параметры под ваши данные (объём, формат, структура, целевые метрики) или нужна настройка окружения в Termux под конкретную задачу, обратитесь в РыбинскЛАБ: поможем с инженерным планом, подбором конфигураций и внедрением стабильного пайплайна.