Пользовательский кластер Databricks, поддерживающий активностьPython

Программы на Python
Anonymous
Пользовательский кластер Databricks, поддерживающий активность

Сообщение Anonymous »

Я обновил свой кластер блоков данных с версии 10.4 LTS до 12.2 LTS, и у меня произошли серьезные изменения в способе использования кластера.
В некотором контексте мы развертываем код Python на Виртуальные машины машинного обучения Azure, которые будут подключаться к кластерам Databricks.
У нас есть один кластер, который используется несколькими алгоритмами. Чтобы все запуски Python могли получить доступ к одному и тому же кластеру и горизонтально масштабироваться. Это полезно для нескольких пользователей, а также снижает наши затраты.
Чтобы еще больше снизить наши затраты, мы решили установить для параметра «Завершение после X минут бездействия» значение 10 минут. Но некоторым нашим алгоритмам необходимо, чтобы сеанс искры оставался активным более 10 минут. Большую часть времени мы записываем DF во временную таблицу или файл, а затем загружаем его. Но иногда лучше оставить сеанс открытым, чтобы поддерживать работоспособность DF.
Для этого мы создали поток, проверяющий кластер каждые 9 минут.

Код: Выделить всё

    if keep_alive:
# Keep Spark alive in a separate thread
self._keep_spark_alive_thread = threading.Thread(
target=self._keep_spark_alive,
daemon=True,
)
self._keep_spark_alive_thread.start()

Код: Выделить всё

    def _keep_spark_alive(self):
"""
Keeps the Spark session alive by running a dummy job and sleeping for a specified interval.

This method runs an infinite loop that periodically executes a dummy job on the Spark session
to prevent it from being terminated due to inactivity. It sleeps for a specified interval between
each execution of the dummy job.

Returns:
None
"""
while True:
# Run a dummy job to keep Spark alive
try:
logging.debug("Keeping Spark alive...")
self._spark.sql(f"SELECT '{self.project_name}'").collect()

# Sleep for 9 minutes
# time.sleep(9 * 60)
time.sleep(20) # 20 seconds for debug
except Exception as e:
logging.debug(f"Error keeping Spark alive: {e}")
При использовании этого потока в кластере 10.4 LTS отправляется Spark sql и поддерживается работоспособность кластера.
Но в 12.2 LTS Spark sql отправляется, НО он игнорируется, как активность, и если в течение 10 минут нет других запросов/действий, кластер отключается.
Я пытался снизить время сна до 20 секунд, и вот что произошло. Каждые 20 секунд я мог видеть активность в журналах Spark, но через 10 минут он начал отключаться. 10 минут 20 секунд, снова запускается из-за треда. Но я теряю сеанс Spark, поскольку кластер перезапускается.
Мои журналы отладки:

Код: Выделить всё

Keeping Spark alive...
Keeping Spark alive...
24/09/30 13:54:05 WARN SparkServiceRPCClient: The cluster seems to be down. A
24/09/30 13:54:06 WARN SparkServiceRPCClient: Cluster xxxx-xxxxxx-xxxxxxxx in
24/09/30 13:54:16 WARN SparkServiceRPCClient: Cluster xxxx-xxxxxx-xxxxxxxx in
24/09/30 13:54:26 WARN SparkServiceRPCClient: Cluster xxxx-xxxxxx-xxxxxxxx in
24/09/30 13:54:37 WARN SparkServiceRPCClient: Cluster xxxx-xxxxxx-xxxxxxxx in
24/09/30 13:54:47 WARN SparkServiceRPCClient: Cluster xxxx-xxxxxx-xxxxxxxx in
24/09/30 13:54:57 WARN SparkServiceRPCClient: Cluster xxxx-xxxxxx-xxxxxxxx in
Error keeping Spark alive: requirement failed: Result for RPC Some(87d40863-d
Keeping Spark alive...
Keeping Spark alive...
Знаете ли вы, почему spark.sql(f"SELECT '{self.project_name}'").collect() не работает как действие? Знает ли новая версия pyspark, что запрос тот же, и поэтому хранится где-то в кеше?

Подробнее здесь: https://stackoverflow.com/questions/790 ... ve-cluster

Вернуться в «Python»