Many companies that want to train models already own years of data in HDFS and Hive tables, a YARN cluster with spare capacity at night, and Spark jobs that know every quirk of that data. Moving it all to a new platform before training anything is slow and risky. The alternative is to use Hadoop for what it is good at, scanning and transforming large data close to where it is stored, and to design the hand-off from that data plane to the GPUs carefully, because that hand-off is where most machine learning pipelines on Hadoop lose their throughput and their reproducibility.

This article walks through that pipeline end to end: building features with Spark on YARN without leaking labels, writing training snapshots in a shape GPUs can consume, scheduling GPUs in YARN, training either on the cluster with Spark's TorchDistributor or off it with a GPU fleet reading HDFS directly, and the failure modes of each. Worked numbers show where the bottleneck really is, and a checklist at the end turns it into next steps.

The pipeline as four hand-offs

Think of the pipeline as four hand-offs. Spark reads raw data, Spark writes a training snapshot, the training job reads its shard of that snapshot, and the result is published with a pointer back to the exact data it saw. Hadoop is strong at the first two: data locality, cheap sequential reads, and a mature engine for joins over billions of rows. It is weak where GPUs are concerned only if you ask it to do things it was never built for, such as serving millions of tiny files or random reads of individual samples.

From data lake to GPU: the four hand-offs that decide throughputRaw HDFS / Hiveevents, logs, tablesSpark on YARNfeature ETL, CPU queueTraining snapshotParquet shards, dated pathManifestschema, counts, splits1 read2 writeA. Train on the clusterTorchDistributor, yarn.io/gpu queueB. Train off the clusterGPU fleet reads via libhdfs3 shard per rankCheckpoints to HDFSsurvive preemptionModel + manifest IDregistry entry4 publishEach numbered arrow is where pipelines usually break: slow scans, tiny files,uneven shards, or a model nobody can trace back to its data.
The training pipeline as four hand-offs. Options A and B differ only in where the GPUs live; the snapshot contract is the same.

The snapshot in the middle is the key design decision. It decouples feature engineering, which is CPU work that YARN schedules well, from training, which is GPU work with very different failure and scaling characteristics. Everything downstream reads an immutable, versioned snapshot, never the live tables.

Features without leakage

Feature engineering for training has one rule that analytics jobs do not: every feature must be computed as of the moment the label was observed. Joining today's customer table to last year's churn labels leaks the future into the past, and the model looks brilliant offline and fails in production. The fix is a point-in-time join: for each labelled event, take the latest feature values with a timestamp at or before the event.

The second rule is deterministic splitting. If training and validation sets are drawn with a random seed per run, rows move between them on every rebuild, and the same user can appear in both. Hash a stable entity key instead, so a given user always lands in the same split:

from pyspark.sql import SparkSession, Window, functions as F

spark = SparkSession.builder.appName("churn-features").enableHiveSupport().getOrCreate()

labels = spark.table("ml.churn_labels")            # user_id, label_ts, churned
feats  = spark.table("warehouse.user_daily_stats") # user_id, stat_ts, sessions_7d, spend_30d, ...

# Point-in-time join: latest feature row at or before each label timestamp.
j = labels.join(feats, "user_id").where(F.col("stat_ts") <= F.col("label_ts"))
w = Window.partitionBy("user_id", "label_ts").orderBy(F.col("stat_ts").desc())
pit = (j.withColumn("rn", F.row_number().over(w)).where("rn = 1").drop("rn"))

# Deterministic split by entity, stable across rebuilds.
bucket = F.abs(F.xxhash64("user_id")) % 100
pit = pit.withColumn("split", F.when(bucket < 80, "train")
                               .when(bucket < 90, "val").otherwise("test"))

snap = "hdfs:///ml/snapshots/churn/2026-10-11"
(pit.repartition(256, "split", "user_id")          # ~256 files per split, not 200,000
    .write.mode("errorifexists")
    .partitionBy("split")
    .option("parquet.block.size", 128 * 1024 * 1024)
    .parquet(snap))

Two details in that write matter for the GPUs. errorifexists makes the snapshot immutable: a dated path is written once and never overwritten. The repartition sets the number of output files so each is a few hundred megabytes with large row groups, which a training reader can stream efficiently. For general feature-engineering patterns in Spark, see Spark ML feature engineering.

Shard files, not sample files

The commonest way to starve a GPU from HDFS is one file per sample: one JPEG per image, one JSON per document. HDFS keeps every file, directory and block as an object in NameNode memory, and a widely used rule of thumb is about 150 bytes of heap per object, so 50 million images cost the NameNode tens of gigabytes before anyone trains anything. Worse, each open is a NameNode round trip plus a DataNode connection, so per-file latency, not bandwidth, sets throughput. The background is in the HDFS small files problem.

Pack samples into large shard files instead. For tabular and text data, Parquet with row groups of 64 to 256 MB is the natural choice, and a training reader iterates row groups. For images, audio and other blobs, store the encoded bytes in a binary column of Parquet, or in a record-oriented format such as TFRecord or WebDataset-style tar shards, with each shard a few hundred megabytes. Shards also make shuffling cheap: shuffle the order of shards per epoch and shuffle samples within a buffer, rather than seeking randomly across the whole dataset.

Worked example: is HDFS fast enough

A team trains an image classifier on 8 GPUs. Each GPU consumes about 1,500 images per second at full utilisation, and images average 120 KB encoded. The data rate the cluster must supply is 8 x 1,500 x 120 KB, about 1.4 GB per second. A modest HDFS cluster with 20 DataNodes, each able to stream a few hundred MB per second from disk, supplies that comfortably when reads are large and sequential.

Now count requests. Stored one image per file, the job opens 12,000 files a second. Each open costs a NameNode RPC and a new DataNode stream, and the data loader spends its time in connection setup; GPUs sit at 30 percent. Packed into 512 MB shards, the same epoch opens about three shards a second. Bandwidth was never the constraint; request rate was. The second check is decode: 12,000 JPEG decodes per second needs many CPU cores, so size data-loader workers accordingly, as discussed in the GPU data loader bottleneck.

Scheduling GPUs in YARN

Since Hadoop 3.1, YARN can schedule GPUs as a first-class resource called yarn.io/gpu, built on the resource-types mechanism explained in YARN resource types. Only Nvidia GPUs are supported; NodeManagers discover devices with nvidia-smi and isolate them with the cgroups devices controller, so a container sees only the GPUs it was given. The scheduler must use the dominant resource calculator, or GPU requests are ignored in placement:

<!-- resource-types.xml -->
<property><name>yarn.resource-types</name><value>yarn.io/gpu</value></property>

<!-- yarn-site.xml on GPU NodeManagers -->
<property><name>yarn.nodemanager.resource-plugins</name><value>yarn.io/gpu</value></property>
<property><name>yarn.nodemanager.resource-plugins.gpu.allowed-gpu-devices</name><value>auto</value></property>
<property><name>yarn.nodemanager.resource-plugins.gpu.path-to-discovery-executables</name>
          <value>/usr/bin/nvidia-smi</value></property>

<!-- capacity-scheduler.xml -->
<property><name>yarn.scheduler.capacity.resource-calculator</name>
  <value>org.apache.hadoop.yarn.util.resource.DominantResourceCalculator</value></property>

Also enable the GPU module in container-executor.cfg. Put GPU hosts behind a node label and a dedicated queue, so CPU-only Spark jobs never occupy them and GPU jobs never wait behind ETL. Spark on YARN translates its own GPU resource requests into yarn.io/gpu:

spark-submit --master yarn --queue gpu \
  --conf spark.yarn.executor.nodeLabelExpression=gpu \
  --conf spark.executor.resource.gpu.amount=1 \
  --conf spark.task.resource.gpu.amount=1 \
  --conf spark.executor.resource.gpu.discoveryScript=/opt/spark/examples/src/main/scripts/getGpusResources.sh \
  --conf spark.executor.cores=8 --conf spark.executor.memory=48g \
  train_churn.py

Training on the cluster with TorchDistributor

With GPUs schedulable, Spark 3.4 and later can launch PyTorch distributed training through pyspark.ml.torch.distributor.TorchDistributor. It runs your training function as a barrier stage, one task per process, sets up the environment torch.distributed needs, pins each task to its GPU, and returns whatever rank 0 returns. Each rank reads its own subset of the snapshot's shard files:

from pyspark.ml.torch.distributor import TorchDistributor

def train(snapshot, epochs):
    import os, torch, torch.distributed as dist
    import pyarrow.dataset as ds, pyarrow.fs as pafs
    dist.init_process_group("nccl")
    rank, world = dist.get_rank(), dist.get_world_size()
    device = torch.device("cuda", 0)          # CUDA_VISIBLE_DEVICES is set per task

    hdfs = pafs.HadoopFileSystem("default")   # uses fs.defaultFS from the Hadoop config
    files = sorted(f.path for f in ds.dataset(f"{snapshot}/split=train",
                   filesystem=hdfs, format="parquet").get_fragments())
    mine = files[rank::world]                 # deterministic shard assignment

    model = torch.nn.parallel.DistributedDataParallel(build_model().to(device))
    opt = torch.optim.AdamW(model.parameters(), lr=1e-3)
    for epoch in range(epochs):
        for batch in iter_batches(hdfs, mine, seed=epoch, batch_size=1024):
            x, y = batch["x"].to(device), batch["y"].to(device)
            loss = torch.nn.functional.binary_cross_entropy_with_logits(model(x), y)
            opt.zero_grad(); loss.backward(); opt.step()
        if rank == 0:
            save_checkpoint(hdfs, model, epoch)   # to hdfs:///ml/checkpoints/...
    dist.destroy_process_group()
    return "done" if rank == 0 else None

result = TorchDistributor(num_processes=4, local_mode=False, use_gpu=True).run(
    train, "hdfs:///ml/snapshots/churn/2026-10-11", 10)

Here build_model, iter_batches and save_checkpoint are your own helpers; iter_batches shuffles shard order by epoch, reads row groups and yields tensors. Barrier mode means all four tasks must be scheduled at once: if the GPU queue can only offer three GPUs, the stage waits rather than starting a partial job. Give rank shards roughly equal sizes, because synchronous training runs at the pace of the slowest rank. Spark-side GPU tuning for ETL is a separate topic, covered in Spark on GPUs.

Training off the cluster

Option B leaves GPUs outside Hadoop: a Kubernetes or Slurm GPU cluster, or cloud GPU instances, reads the snapshot directly. PyArrow's HadoopFileSystem wraps libhdfs, so the training hosts need a JVM, JAVA_HOME, libhdfs.so and a CLASSPATH built from hadoop classpath --glob, plus network reach to every DataNode, not just the NameNode. On a Kerberos cluster they also need credentials. Alternatively, a scheduled DistCp job copies each finished snapshot to object storage the GPU fleet already reads; the copy costs time and storage but removes the Hadoop client from GPU hosts entirely and lets the cluster stay private.

Choose by where the GPUs are and how often data changes. On-cluster training reuses the queue, security and locality you already run; off-cluster training lets a GPU platform team operate drivers, NCCL and images without touching Hadoop. Both read the same immutable snapshot, which is what keeps them interchangeable.

Failure modes

  • Idle GPUs. Utilisation below 50 percent with data-loader workers pegged. Cause: small files, too few loader workers, or decode on the main process. Fix: shard, add workers, prefetch.
  • Barrier stage never starts. TorchDistributor waits indefinitely because the queue cannot grant every GPU at once, often because ETL containers hold GPU hosts. Fix: GPU node label, dedicated queue, and a job-level timeout.
  • Preempted training. YARN preemption or a NodeManager loss kills one rank and the whole synchronous job fails. Fix: checkpoint to HDFS every N minutes and resume; make the GPU queue non-preemptable if jobs are long.
  • Token expiry. On secure clusters HDFS delegation tokens are renewed daily and expire after seven days by default, so a long training run or a long-lived external reader fails mid-epoch. Fix: keytab-based login on long-running readers, and keep jobs inside the token lifetime.
  • Leakage and drift between runs. A model improves suspiciously after a rebuild. Cause: non-point-in-time joins or random splits. Fix: point-in-time joins, hashed splits, and a manifest recording row counts and the schema per snapshot.
  • Unreproducible models. Nobody can say which data trained the production model. Fix: register the snapshot path and manifest hash with every model, and use HDFS snapshots or immutable paths so the data cannot change underneath it.

Trade-offs

Using Hadoop for AI is a good trade when data is already there, volumes are large, and the cluster has spare capacity. It is a poor one when the team must build GPU operations inside YARN from scratch for one model, or when training data is small enough to fit on one machine, where a single-node job reading one Parquet export beats any distributed design. Lakehouse platforms with object storage and table formats are increasingly where new pipelines start; the snapshot pattern here ports to them unchanged, so building it on Hadoop today is not wasted work. For serving the same features online, see feature stores.

What to do next

  1. Pick one model and write its features as a point-in-time join with a hashed train, validation and test split.
  2. Write training data as immutable, dated Parquet snapshots of a few hundred megabytes per file, plus a manifest.
  3. Measure GPU utilisation and data-loader throughput on one GPU before scaling out.
  4. If you train on the cluster, enable yarn.io/gpu, add a GPU node label and a dedicated queue.
  5. Run a four-process TorchDistributor job with rank-based shard assignment and checkpoints to HDFS.
  6. If you train off the cluster, test libhdfs access and Kerberos from one GPU host, or schedule DistCp exports.
  7. Record the snapshot path and manifest hash with every registered model.
Key takeaway: Hadoop can feed machine learning well if you treat the hand-off as a contract: build features with point-in-time joins and hashed splits, write immutable Parquet snapshots in large shards, schedule GPUs through yarn.io/gpu or read snapshots from a separate GPU fleet, and record which snapshot trained every model.