【问题标题】:Airflow/Composer UnzipOperator CrashesAirflow/Composer UnzipOperator 崩溃
【发布时间】:2020-10-21 20:07:00
【问题描述】:

我尝试运行我的 DAG 来处理数 GB 的压缩文件 (zip)。 dag 使用以下运算符:

    ... # Other operators

    >> UnzipOperator(task_id="unzip_archive",
    path_to_zip_file=archive_path,
    path_to_unzip_contents=unzipped_f_path)

    ... # Other operators

DAG 似乎崩溃了,有这个日志:

-------------------------------------------------------------------------------
Starting attempt 1 of 
-------------------------------------------------------------------------------

[2020-06-26 15:10:56,157] {models.py:1599} INFO - Executing <Task(UnzipOperator): unzip_archive> on 2020-06-26T14:34:15+00:00
[2020-06-26 15:10:56,163] {base_task_runner.py:118} INFO - Running: ['bash', '-c', 'airflow run fota_integration.fota_user_profile unzip_archive 2020-06-26T14:34:15+00:00 --job_id 227634 --pool integration --raw -sd DAGS_FOLDER/fota/fota_integration.py --cfg_path /tmp/tmpdht77jvb']
[2020-06-26 15:11:12,208] {base_task_runner.py:101} INFO - Job 227634: Subtask unzip_archive [2020-06-26 15:11:12,207] {settings.py:176} INFO - settings.configure_orm(): Using pool settings. pool_size=5, pool_recycle=1800, pid=297
[2020-06-26 15:11:21,582] {base_task_runner.py:101} INFO - Job 227634: Subtask unzip_archive [2020-06-26 15:11:21,537] {default_celery.py:90} WARNING - You have configured a result_backend of redis://airflow-redis-service.default.svc.cluster.local:6379/0, it is highly recommended to use an alternative result_backend (i.e. a database).
[2020-06-26 15:11:21,692] {base_task_runner.py:101} INFO - Job 227634: Subtask unzip_archive [2020-06-26 15:11:21,692] {__init__.py:51} INFO - Using executor CeleryExecutor
[2020-06-26 15:11:28,126] {base_task_runner.py:101} INFO - Job 227634: Subtask unzip_archive [2020-06-26 15:11:28,022] {app.py:52} WARNING - Using default Composer Environment Variables. Overrides have not been applied.
[2020-06-26 15:11:28,650] {base_task_runner.py:101} INFO - Job 227634: Subtask unzip_archive [2020-06-26 15:11:28,617] {configuration.py:522} INFO - Reading the config from /etc/airflow/airflow.cfg
[2020-06-26 15:11:29,279] {base_task_runner.py:101} INFO - Job 227634: Subtask unzip_archive [2020-06-26 15:11:29,270] {configuration.py:522} INFO - Reading the config from /etc/airflow/airflow.cfg
[2020-06-26 15:11:31,908] {base_task_runner.py:101} INFO - Job 227634: Subtask unzip_archive [2020-06-26 15:11:31,900] {models.py:273} INFO - Filling up the DagBag from /home/airflow/gcs/dags/fota/fota_integration.py
[2020-06-26 15:11:35,673] {base_task_runner.py:101} INFO - Job 227634: Subtask unzip_archive [2020-06-26 15:11:35,670] {cli.py:520} INFO - Running <TaskInstance: fota_integration.fota_user_profile.unzip_archive 2020-06-26T14:34:15+00:00 [running]> on host airflow-worker-778c879665-zgjbn
[2020-06-26 15:11:36,086] {base_task_runner.py:101} INFO - Job 227634: Subtask unzip_archive [2020-06-26 15:11:36,080] {zip.py:136} INFO - Executing UnzipOperator.execute(context)
[2020-06-26 15:11:36,086] {base_task_runner.py:101} INFO - Job 227634: Subtask unzip_archive [2020-06-26 15:11:36,084] {zip.py:138} INFO - path_to_zip_file: /home/airflow/gcs/data/zips/20200626_USER_PROFILE.zip
[2020-06-26 15:11:36,086] {base_task_runner.py:101} INFO - Job 227634: Subtask unzip_archive [2020-06-26 15:11:36,084] {zip.py:140} INFO - path_to_unzip_contents: /home/airflow/gcs/data/raw_data/20200626/USER_PROFILE/
[2020-06-26 15:11:36,365] {base_task_runner.py:101} INFO - Job 227634: Subtask unzip_archive [2020-06-26 15:11:36,365] {zip.py:155} INFO - Created zip file object '<zipfile.ZipFile filename='/home/airflow/gcs/data/zips/20200626_USER_PROFILE.zip' mode='r'>' from path '/home/airflow/gcs/data/zips/20200626_USER_PROFILE.zip'
[2020-06-26 15:11:36,365] {base_task_runner.py:101} INFO - Job 227634: Subtask unzip_archive [2020-06-26 15:11:36,365] {zip.py:158} INFO - Extracting all the contents to '/home/airflow/gcs/data/raw_data/20200626/USER_PROFILE/'
[2020-06-26 15:43:25,451] {helpers.py:250} INFO - Sending Signals.SIGTERM to GPID 297
[2020-06-26 15:43:25,529] {models.py:1641} ERROR - Received SIGTERM. Terminating subprocesses.
[2020-06-26 15:43:25,533] {base_task_runner.py:101} INFO - Job 227634: Subtask unzip_archive [2020-06-26 15:43:25,529] {models.py:1641} ERROR - Received SIGTERM. Terminating subprocesses.
[2020-06-26 15:43:29,047] {helpers.py:250} INFO - Sending Signals.SIGTERM to GPID 297
[2020-06-26 15:43:29,048] {models.py:1641} ERROR - Received SIGTERM. Terminating subprocesses.
[2020-06-26 15:43:29,054] {base_task_runner.py:101} INFO - Job 227634: Subtask unzip_archive [2020-06-26 15:43:29,048] {models.py:1641} ERROR - Received SIGTERM. Terminating subprocesses.
[2020-06-26 15:43:38,943] {helpers.py:232} INFO - Process psutil.Process(pid=297, status='terminated') (297) terminated with exit code 0

然后当我尝试重新运行时,我得到了这个日志:

-------------------------------------------------------------------------------
Starting attempt 2 of 
-------------------------------------------------------------------------------

[2020-07-01 09:14:33,615] {models.py:1599} INFO - Executing <Task(UnzipOperator): unzip_archive> on 2020-06-26T14:34:15+00:00
[2020-07-01 09:14:33,616] {base_task_runner.py:118} INFO - Running: ['bash', '-c', 'airflow run fota_integration.fota_user_profile unzip_archive 2020-06-26T14:34:15+00:00 --job_id 231875 --pool integration --raw -sd DAGS_FOLDER/fota/fota_integration.py --cfg_path /tmp/tmp0uwqdre9']

我想说的是,显然内存不足会触发 SIGKILL ......但是当我检查UnzipOperator 库的代码时,我发现有效负载没有导入到内存中。

这是解压缩操作符的主循环。

    with self.open(member, pwd=pwd) as source, \
         open(targetpath, "wb") as target:
        shutil.copyfileobj(source, target)

位于_extract_member(self, member, targetpath, pwd),用于extractall(self, path=None, members=None, pwd=None),用于UnzipOperator类的execute(self, context)

我的问题:

  1. 为什么会发生 SIGKILL?
  2. 我该如何解决这个问题?

编辑:

我将另一个问题与这个问题联系起来,因为我不确定其中分析的有效性。它与无法编写 python 文件而不将其完全加载到内存中有关。 Find it here.

【问题讨论】:

  • SIGTERM 通常是作业心跳失败的结果,这可能是过载的间接症状。您可以在此期间发布 Airflow 工作人员的日志吗?
  • @hexacyanide,不知道该怎么做。我正在使用作曲家(kubernetes 集群)
  • 为了帮助隔离 SIGTERM 的原因,您可以在其中一个工作容器中安装 strace,然后尝试跟踪导致 SIGTERM 的进程的 PID/名称。这是你可以做到的cloud.google.com/container-optimized-os/docs/how-to/…
  • 要进一步解决此问题,您可以在cloud.google.com/support/docs/issue-trackers 上报告它

标签: python google-cloud-platform zip airflow google-cloud-composer


【解决方案1】:

这是一个不幸的案例。通货膨胀的错误没有通过 Python 异常链提出。

你能进入实例直接解压缩对象吗?希望这会产生足够的错误数据供您采取行动。或者,如果您找不到问题,您可以使用 BashOperator 之类的工具。


我可能知道发生了什么 [2020-06-26 15:11:36,365] --- [2020-06-26 15:43:25,451]...这是 30 分钟...这是预期的吗?... ..我的猜测是解压缩操作员根本没有向执行程序或调度程序提供足够的数据,可能调度程序本身认为工作已完成或超时......一些超时逻辑被击中并且进程关闭。

有一些方法可以更新或延长此超时。希望这足以提供帮助

【讨论】:

  • 确实,一个 bash 操作符可以做到这一点......除了它让我不知道它到底在哪里失败。我需要了解
  • 我不确定是超时,因为文件是 5 Gb,再加上简单的重新运行至少会产生相同的响应,但第二次运行是空的......进程突然中断
猜你喜欢
  • 2019-05-08
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2020-07-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多