【问题标题】:Using spark_sklearn with a kubernetes cluster将 spark_sklearn 与 Kubernetes 集群一起使用
【发布时间】:2019-11-07 23:21:16
【问题描述】:

我正在做一个机器学习项目。我最初使用 scikit-learn (sklearn) 库。在模型优化过程中,我使用了 sklearn 中的经典 GridSearchCV 类。它目前使用运行 api 的主机(joblib 库)中的所有资源进行并行化。下面你有一个例子,

from sklearn                  import datasets
from sklearn.ensemble         import RandomForestClassifier
from sklearn.model_selection  import GridSearchCV
from datetime                 import datetime
import numpy as np

def count_trials(param_grid):
    total_trials = 0
    for k,v in param_grid.items():
        total_trials += len(v)

    return total_trials

# Load data
digits = datasets.load_digits()

X, y = digits.data, digits.target

print("")
print("Iris data set: ")
print("X:      {}".format(X.shape))
print("labels: {}".format(np.unique(y)))
print("")

param_grid = {"max_depth":         [3, None],
              "max_features":      ["auto"],
              "min_samples_split": [2, 3,4,5,10,20],
              "min_samples_leaf":  [1, 3,4,6,10,20],
              "bootstrap":         [True],
              "criterion":         ["entropy"],
              "n_estimators":      [40, 80],
              }
cv = 5

n_models = count_trials(param_grid)

print("trying {} models, with CV = {}. A total of {} fits.".format(n_models,cv,n_models*cv))

start_time = datetime.now()
print("Starting at {}".format(start_time))
gs = GridSearchCV(estimator = RandomForestClassifier(),
                  param_grid=param_grid,
                  cv=cv,
                  refit=True,
                  scoring="accuracy",
                  n_jobs = -1)
gs.fit(X, y)
end_time   = datetime.now()

print("Ending   at {}".format(end_time))
print("\n total time = {}\n".format(end_time - start_time))

我最近发现它已被扩展为使用 spark 集群的资源(pyspark 和 spark_sklearn 库)。我设法设置了一个包含一个主节点和两个工作节点的 spark 集群。下面的代码运行与以前相同的任务,但使用 spark 集群资源。

from sklearn                  import datasets
from sklearn.ensemble         import RandomForestClassifier
from sklearn.model_selection  import GridSearchCV as SKGridSearchCV
from spark_sklearn            import GridSearchCV as SparkGridSearchCV
from pyspark                  import SparkConf, SparkContext
from datetime                 import datetime
import numpy as np

def get_context():
    sc_conf = SparkConf()
    sc_conf.setAppName("test-sklearn-spark-app")
    sc_conf.setMaster('spark://<master-IP>:7077')
    sc_conf.set('spark.cores.max', '40')
    sc_conf.set('spark.logConf', True)
    print(sc_conf.getAll())

    return SparkContext(conf=sc_conf)

def count_trials(param_grid):
    total_trials = 0
    for k,v in param_grid.items():
        total_trials += len(v)

    return total_trials

# Load data
digits = datasets.load_digits()

X, y = digits.data, digits.target

print("")
print("Iris data set: ")
print("X:      {}".format(X.shape))
print("labels: {}".format(np.unique(y)))
print("")

param_grid = {"max_depth":         [3, None],
              "max_features":      ["auto"],
              "min_samples_split": [2, 3,4,5,10,20],
              "min_samples_leaf":  [1, 3,4,6,10,20],
              "bootstrap":         [True],
              "criterion":         ["entropy"],
              "n_estimators":      [40, 80],
              }
cv = 5

n_models = count_trials(param_grid)

print("trying {} models, with CV = {}. A total of {} fits.".format(n_models,cv,n_models*cv))

sc = get_context()
gs = SparkGridSearchCV(sc = sc,
                       estimator = RandomForestClassifier(),
                       param_grid=param_grid,
                       cv=cv,
                       refit=True,
                       scoring="accuracy",
                       n_jobs = -1)

start_time = datetime.now()
print("Starting at {}".format(start_time))
gs.fit(X, y)
end_time   = datetime.now()

print("Ending   at {}".format(end_time))
print("\n total time = {}\n".format(end_time - start_time))

其中 master-IP 是主节点的 IP。该代码完美运行,使用了 spark 集群中的所有可用资源。

然后我配置了一个包含一个主节点和一个从节点的 Kubernetes 集群。然后我运行与以前相同的代码,但用

替换该行
sc_conf.setMaster('spark://<master-IP>:7077')

通过

sc_conf.setMaster('k8s://<master-IP>:<PORT>')

其中master-IP和PORT是我在master节点上运行命令得到的,

kubectl  cluster-info

问题是我的代码不再工作了。它显示以下错误消息,

19/11/07 12:57:32 ERROR Utils: Uncaught exception in thread kubernetes-executor-snapshots-subscribers-1
org.apache.spark.SparkException: Must specify the executor container image
    at org.apache.spark.deploy.k8s.features.BasicExecutorFeatureStep$$anonfun$5.apply(BasicExecutorFeatureStep.scala:40)
    at org.apache.spark.deploy.k8s.features.BasicExecutorFeatureStep$$anonfun$5.apply(BasicExecutorFeatureStep.scala:40)
    at scala.Option.getOrElse(Option.scala:121)
    at org.apache.spark.deploy.k8s.features.BasicExecutorFeatureStep.<init>(BasicExecutorFeatureStep.scala:40)
    at org.apache.spark.scheduler.cluster.k8s.KubernetesExecutorBuilder$$anonfun$$lessinit$greater$default$1$1.apply(KubernetesExecutorBuilder.scala:26)
    at org.apache.spark.scheduler.cluster.k8s.KubernetesExecutorBuilder$$anonfun$$lessinit$greater$default$1$1.apply(KubernetesExecutorBuilder.scala:26)
    at org.apache.spark.scheduler.cluster.k8s.KubernetesExecutorBuilder.buildFromFeatures(KubernetesExecutorBuilder.scala:43)
    at org.apache.spark.scheduler.cluster.k8s.ExecutorPodsAllocator$$anonfun$org$apache$spark$scheduler$cluster$k8s$ExecutorPodsAllocator$$onNewSnapshots$1.apply$mcVI$sp(ExecutorPodsAllocator.scala:133)
    at scala.collection.immutable.Range.foreach$mVc$sp(Range.scala:160)
    at org.apache.spark.scheduler.cluster.k8s.ExecutorPodsAllocator.org$apache$spark$scheduler$cluster$k8s$ExecutorPodsAllocator$$onNewSnapshots(ExecutorPodsAllocator.scala:126)
    at org.apache.spark.scheduler.cluster.k8s.ExecutorPodsAllocator$$anonfun$start$1.apply(ExecutorPodsAllocator.scala:68)
    at org.apache.spark.scheduler.cluster.k8s.ExecutorPodsAllocator$$anonfun$start$1.apply(ExecutorPodsAllocator.scala:68)
    at org.apache.spark.scheduler.cluster.k8s.ExecutorPodsSnapshotsStoreImpl$$anonfun$org$apache$spark$scheduler$cluster$k8s$ExecutorPodsSnapshotsStoreImpl$$callSubscriber$1.apply$mcV$sp(ExecutorPodsSnapshotsStoreImpl.scala:102)
    at org.apache.spark.util.Utils$.tryLogNonFatalError(Utils.scala:1340)
    at org.apache.spark.scheduler.cluster.k8s.ExecutorPodsSnapshotsStoreImpl.org$apache$spark$scheduler$cluster$k8s$ExecutorPodsSnapshotsStoreImpl$$callSubscriber(ExecutorPodsSnapshotsStoreImpl.scala:99)
    at org.apache.spark.scheduler.cluster.k8s.ExecutorPodsSnapshotsStoreImpl$$anonfun$addSubscriber$1.apply$mcV$sp(ExecutorPodsSnapshotsStoreImpl.scala:71)
    at org.apache.spark.scheduler.cluster.k8s.ExecutorPodsSnapshotsStoreImpl$$anon$1.run(ExecutorPodsSnapshotsStoreImpl.scala:107)
    at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511)
    at java.util.concurrent.FutureTask.runAndReset(FutureTask.java:308)
    at java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.access$301(ScheduledThreadPoolExecutor.java:180)
    at java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.run(ScheduledThreadPoolExecutor.java:294)
    at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
    at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
    at java.lang.Thread.run(Thread.java:748)

好像说我要指定一个docker镜像,但是我不知道该怎么做。

有人有这方面的经验吗?我一直在网上四处寻找,但没有答案。

提前谢谢你,

麦芽酒

【问题讨论】:

    标签: apache-spark kubernetes scikit-learn pyspark


    【解决方案1】:

    我建议你先阅读doc

    当您在 Kubernetes 上运行 Spark 时,您的进程将提交到 Docker 容器内的 Kubernetes 集群,在提交作业时应该知道这一点。设置为 SparkConf:

    • spark.kubernetes.container.image

    • spark.kubernetes.driver.container.image
    • spark.kubernetes.executor.container.image

    除了确保您的 Kubernetes 集群可以拉取这些镜像外,最简单的方法是将它们推送到 DockerHub。关注guide,了解如何构建 Spark Docker 镜像。

    您似乎在客户端模式下运行您的作业,因此请考虑networking notes。基本上,您需要确保您的 Driver 进程(可能在您的本地机器上运行)可以从 Kubernetes 网络(特别是执行程序 Pod)中访问,这可能不是那么明显。同样,您的 Driver 进程应该具有对 Executor Pod 的网络访问权限。实际上,从本地工作站以客户端模式将 Spark Jobs 提交到远程 Kubernetes 集群确实很棘手,我建议您先尝试使用集群模式。

    如果您想在集群模式下提交您的作业,您需要确保您的作业工件(在您的情况下是 python 脚本)及其依赖项可以从 Spark Driver 和 Executor Pods 访问(最简单的方法是将您的脚本包含 Spark 类路径上 Spark Docker 映像中的所有依赖项)。

    而不是它应该以与通常相同的方式为您工作。

    您还可以参考 Helm chart of Spark on Kubernetes cluster,其中包括 Jupyter 笔记本集成,它可以更轻松地在 Kubernetes 上运行交互式 Spark 会话。

    【讨论】:

      猜你喜欢
      • 2016-02-12
      • 1970-01-01
      • 2012-01-13
      • 2021-04-10
      • 1970-01-01
      • 2017-12-17
      • 2019-03-25
      • 2016-04-12
      • 2019-05-28
      相关资源
      最近更新 更多