diff --git a/datasets/add_months.sh b/datasets/add_months.sh index 8f331c7..3c745dc 100755 --- a/datasets/add_months.sh +++ b/datasets/add_months.sh @@ -109,15 +109,31 @@ parallel --joblog add_months_joblog.txt --results add_months_logs \ source "$SPARK_CONF_DIR/spark-env.sh" export PYSPARK_PYTHON="$PYTHON" + +# Compute executor resources from what spark-env.sh detected. +# One executor gets all worker cores and ~85% of worker memory (leaving +# headroom for the driver and overhead). +_EXEC_CORES=${SPARK_WORKER_CORES:-$(nproc)} +_WORKER_MEM_GB=${SPARK_WORKER_MEMORY%g} +_EXEC_MEM_GB=$(( _WORKER_MEM_GB * 85 / 100 )) +_DRIVER_MEM_GB=$(( _WORKER_MEM_GB * 10 / 100 )) +SPARK_SUBMIT_CONF=( + --conf "spark.executor.cores=${_EXEC_CORES}" + --conf "spark.executor.memory=${_EXEC_MEM_GB}g" + --conf "spark.driver.memory=${_DRIVER_MEM_GB}g" +) + start_spark_cluster.sh spark-submit --master "spark://$(hostname):$SPARK_MASTER_PORT" \ + "${SPARK_SUBMIT_CONF[@]}" \ comments_part2.py \ --indir="$TEMP_COMMENTS" \ --out_by_subreddit="$STAGING_COMMENTS_SUB" \ --out_by_author="$STAGING_COMMENTS_AUTH" spark-submit --master "spark://$(hostname):$SPARK_MASTER_PORT" \ + "${SPARK_SUBMIT_CONF[@]}" \ submissions_part2.py \ --indir="$TEMP_SUBMISSIONS" \ --out_by_subreddit="$STAGING_SUBMISSIONS_SUB" \