【发布时间】:2020-08-09 17:56:15
【问题描述】:
我试图将经过训练的 Faiss 索引部署到 PySpark 并进行分布式搜索。所以整个过程包括:
- 预处理
- 加载 Faiss 索引(~15G)并进行 Faiss 搜索
- 后处理并写入 HDFS
我将每个任务的 CPU 设置为 10 (spark.task.cpus=10) 以便进行多线程搜索。但是步骤 1 和步骤 3 每个任务只能使用 1 个 CPU。为了利用所有 CPU,我想在第 1 步和第 3 步之前设置spark.task.cpus=1。我尝试了RuntimeConfig 的设置方法,但它似乎让我的程序卡住了。关于如何在运行时更改配置或如何优化此问题的任何建议?
代码示例:
def load_and_search(x, model_path):
faiss_idx = faiss.read_index(model_path)
q_vec = np.concatenate(x)
_, idx_array = faiss_idx.search(q_vec, k=10)
return idx_array
data = sc.textFile(input_path)
# preprocess, only used one cpu per task
data = data.map(lambda x: x)
# load faiss index and search, used multiple cpus per task
data = data.mapPartitioins(lambda x: load_and_search(x, model_path))
# postprocess and write, one cpu per task
data = data.map(lambda x: x).saveAsTextFile(result_path)
【问题讨论】:
-
您能否给出一个更具体的代码示例来说明您正在尝试做的事情,以便更容易帮助您进行优化?
-
@user3689574 感谢您的评论。已添加示例。
-
load_and_search这样做不是 Spark 代码吗?
标签: apache-spark pyspark