Domino Distributed Computing

作者 dominodatalabd86698d74d56无许可证7 个星标收录于 2026年10月8日更新于 2026年10月8日仓库2天前更新

Work with distributed computing frameworks in Domino including Apache Spark, Ray, and Dask clusters. Covers cluster configuration, on-demand clusters, choosing between frameworks, PySpark usage, and scaling workloads. Use when processing large datasets, parallel ML training, or running distributed compute jobs.

AI 生成的概览

指导在 Domino 中运行 Spark、Ray 和 Dask 分布式集群,涵盖配置、框架选择与扩缩容。

功能
该技能提供在 Domino 中使用分布式计算框架的说明,涵盖 Apache Spark、Ray、Dask 和 MPI。它讲解如何通过 Domino 界面或 Python SDK 启动按需集群、连接各框架,并运行数据处理、分布式训练、超参数调优和 GPU 工作负载。它还提供框架选择、自动扩缩容、最佳实践和故障排查方面的指导。
适用场景
适用于处理大型数据集、扩展数据处理或机器学习训练、配置分布式集群设置,或判断哪种框架适合某项工作负载。它面向可使用 Spark、Ray 或 Dask 集群的 Domino 环境。
运行要求
需要可访问 Spark、Ray 或 Dask 集群的 Domino 环境,以及相关 Python 库(例如 pyspark、ray、dask、dask-ml)。以编程方式启动集群需要 Domino Python SDK。不包含脚本,仅为说明文档。

Domino Distributed Computing Skill

Description

This skill helps users work with distributed computing frameworks in Domino - Spark, Ray, and Dask clusters for scaling compute-intensive workloads.

Activation

Activate this skill when users want to:

  • Run Spark, Ray, or Dask clusters in Domino
  • Scale data processing or ML training
  • Configure distributed cluster settings
  • Understand when to use each framework

Supported Frameworks

FrameworkBest For
Apache SparkLarge-scale data processing, SQL, ETL
RayDistributed ML, hyperparameter tuning, RL
DaskParallel pandas, NumPy at scale
MPIScientific computing, HPC workloads

When to Use Each Framework

Spark

  • Processing terabyte-scale data
  • SQL analytics on big data
  • ETL pipelines
  • Structured data processing

Ray

  • Distributed model training
  • Hyperparameter optimization
  • Reinforcement learning
  • Generic Python parallelization

Dask

  • Scaling pandas workflows
  • Parallel NumPy operations
  • Lazy evaluation needed
  • Familiar pandas/NumPy API preferred

Launching On-Demand Clusters

Via Domino UI

  1. Start a workspace or job
  2. Check Attach compute cluster
  3. Select:
    • Cluster Type: Spark, Ray, or Dask
    • Worker Count: Number of workers
    • Hardware Tier: Resources per worker
    • Auto-scaling: Enable/disable
  4. Launch

Via Python SDK

python
from domino import Domino
domino = Domino("project-owner/project-name")
# Start workspace with Spark clusterworkspace = domino.workspace_start(    hardware_tier_name="medium",    cluster_config={        "clusterType": "Spark",        "workerCount": 4,        "workerHardwareTier": "medium",        "masterHardwareTier": "medium"    })

Apache Spark

Connecting to Spark

python
from pyspark.sql import SparkSession
# Domino auto-configures Sparkspark = SparkSession.builder.getOrCreate()
# Check configurationprint(f"Spark version: {spark.version}")print(f"Executors: {spark.sparkContext.defaultParallelism}")

Reading Data

python
# Read CSVdf = spark.read.csv("/mnt/data/dataset/data.csv", header=True, inferSchema=True)
# Read Parquetdf = spark.read.parquet("/mnt/data/dataset/")
# Read from databasedf = spark.read.jdbc(    url="jdbc:postgresql://host:5432/db",    table="schema.table",    properties={"user": "user", "password": "pass"})

Processing Data

python
from pyspark.sql import functions as F
# Transformationsresult = df.filter(F.col("value") > 100) \    .groupBy("category") \    .agg(F.mean("value").alias("avg_value")) \    .orderBy("avg_value", ascending=False)
# Show resultsresult.show()

Machine Learning with Spark MLlib

python
from pyspark.ml.feature import VectorAssemblerfrom pyspark.ml.classification import RandomForestClassifierfrom pyspark.ml import Pipeline
# Prepare featuresassembler = VectorAssembler(    inputCols=["feature1", "feature2", "feature3"],    outputCol="features")
# Create modelrf = RandomForestClassifier(    featuresCol="features",    labelCol="label",    numTrees=100)
# Build pipelinepipeline = Pipeline(stages=[assembler, rf])model = pipeline.fit(train_df)predictions = model.transform(test_df)

Writing Results

python
# Write Parquet (recommended)result.write.parquet("/mnt/artifacts/output/", mode="overwrite")
# Write CSVresult.write.csv("/mnt/artifacts/output.csv", header=True)

Ray

Connecting to Ray

python
import ray
# Domino auto-initializes Ray# Or manually connectray.init(address="auto")
print(f"Cluster resources: {ray.cluster_resources()}")

Parallel Tasks

python
import ray
@ray.remotedef process_item(item):    # Your processing logic    return item * 2
# Run in parallelitems = [1, 2, 3, 4, 5]futures = [process_item.remote(item) for item in items]results = ray.get(futures)print(results)  # [2, 4, 6, 8, 10]

Distributed Training with Ray Train

python
from ray import trainfrom ray.train import ScalingConfigfrom ray.train.torch import TorchTrainer
def train_func():    # Training logic    model = create_model()    for epoch in range(10):        train_epoch(model)        train.report({"loss": loss})
trainer = TorchTrainer(    train_func,    scaling_config=ScalingConfig(num_workers=4, use_gpu=True))result = trainer.fit()

Hyperparameter Tuning with Ray Tune

python
from ray import tunefrom ray.tune import CLIReporter
def objective(config):    # Training with hyperparameters    model = train_model(        learning_rate=config["lr"],        batch_size=config["batch_size"]    )    return {"accuracy": accuracy}
analysis = tune.run(    objective,    config={        "lr": tune.loguniform(1e-4, 1e-1),        "batch_size": tune.choice([32, 64, 128])    },    num_samples=20,    progress_reporter=CLIReporter())
print(f"Best config: {analysis.best_config}")

Dask

Connecting to Dask

python
from dask.distributed import Client
# Domino auto-configures Daskclient = Client()
print(f"Dashboard: {client.dashboard_link}")print(f"Workers: {len(client.scheduler_info()['workers'])}")

Dask DataFrames (Parallel pandas)

python
import dask.dataframe as dd
# Read large CSV filesdf = dd.read_csv("/mnt/data/dataset/*.csv")
# Parallel operations (lazy)result = df.groupby("category")["value"].mean()
# Executecomputed_result = result.compute()

Dask Arrays (Parallel NumPy)

python
import dask.array as da
# Create large arrayx = da.random.random((100000, 100000), chunks=(1000, 1000))
# Operations (lazy)result = x.mean()
# Computevalue = result.compute()

Dask ML

python
from dask_ml.model_selection import GridSearchCVfrom sklearn.ensemble import RandomForestClassifier
# Distributed hyperparameter searchparam_grid = {    "n_estimators": [100, 200, 300],    "max_depth": [10, 20, 30]}
grid_search = GridSearchCV(    RandomForestClassifier(),    param_grid,    cv=3)
grid_search.fit(X_train, y_train)print(f"Best params: {grid_search.best_params_}")

GPU Clusters

Spark RAPIDS

python
# Use GPU-accelerated Sparkspark = SparkSession.builder \    .config("spark.rapids.sql.enabled", "true") \    .getOrCreate()
# Operations automatically use GPUdf = spark.read.parquet("/mnt/data/large_dataset/")result = df.groupBy("category").agg({"value": "mean"})

Ray with GPUs

python
@ray.remote(num_gpus=1)def train_on_gpu():    import torch    device = torch.device("cuda")    # GPU training logic    return model
# Run on GPU workersfutures = [train_on_gpu.remote() for _ in range(4)]

Autoscaling

Enable Autoscaling

Configure clusters to scale based on workload:

python
cluster_config = {    "clusterType": "Spark",    "workerCount": 2,    "maxWorkerCount": 10,  # Scale up to 10    "autoScaling": True}

Monitor Scaling

View cluster status in Domino UI or via dashboard URLs.

Best Practices

1. Choose Right Framework

  • SQL/ETL: Spark
  • ML/Parallel Python: Ray
  • Pandas at scale: Dask

2. Right-size Clusters

  • Start small, scale up
  • Monitor resource usage
  • Use autoscaling when unsure

3. Data Locality

python
# Keep data close to compute# Use Domino Datasets or cloud storage in same regiondf = spark.read.parquet("/mnt/data/dataset/")

4. Persist Intermediate Results

python
# Cache frequently used DataFramesdf.cache()df.persist()

Troubleshooting

Cluster Won't Start

  • Check hardware tier availability
  • Verify cluster environment builds
  • Review cluster logs

Out of Memory

  • Increase worker memory
  • Add more workers
  • Optimize code (reduce shuffles)

Slow Performance

  • Check data locality
  • Review partition sizes
  • Monitor cluster dashboard

Documentation Reference

来源与署名

来源:dominodatalab/domino-claude-plugin位于skills/distributed-computing提交d86698d

许可证: 无许可证

内容归原作者所有。SourceWeft 从公开仓库中收录这些内容。

举报或申请下架