【问题标题】:How do I emit Airflow logs from a DataProcPySparkOperator如何从 DataProcPySparkOperator 发出气流日志
【发布时间】:2020-12-08 21:16:01
【问题描述】:

我已经看到了一些关于从 PythonOperatorconfiguring Airflow logs 发出日志的信息,但还没有发现任何可以让我从容器化进程中发出日志的东西,例如DataProcPySparkOperator。

我已经在 Operator 集群中运行的 pyspark 脚本的顶部添加了以下内容:

import logging
logging.info('Test bare logger')
for ls in ['airflow', 'airflow.task', __name__]:
    l = logging.getLogger(ls)
    l.info('Test {} logger'.format(ls))
print('Test print() logging')

它不会产生任何输出,尽管 Operator 脚本会按预期运行。

我假设我可以从集群内建立到云存储(或数据库)的连接,也许可以捎带用于读取和写入文件的现有连接,但是......这似乎就像为了一个共同的需要做很多工作。我非常希望在计算的中间阶段获得关于记录数量或其他数据的偶尔引用的状态检查。

Airflow 是否默认在集群中设置 Python 记录器?如果是,我该如何访问它?

【问题讨论】:

  • 此脚本是在 GCP 还是本地集群上运行,您要将日志发送到哪里?将一般 python 日志导出到 GCP 日志浏览器的一种可能方法是 - cloud.google.com/logging/docs/setup/python
  • 脚本正在 GCP 上运行。感谢您的链接。会看的。如果可能的话,我希望日志包含在 Airflow 的正常任务执行日志中。

标签: python logging airflow google-cloud-dataproc


【解决方案1】:

如果您使用 Cloud Composer 运行 Airflow,您应该知道 Cloud Composer 仅包含 Airflow 日志和 Streaming 日志。

当您在 Dataproc 集群中运行 PySpark 作业时,作业驱动程序输出存储在 Cloud Storage 中(请参阅Accessing job driver output

您还可以启用 Dataproc 以将作业驱动程序日志保存在 Cloud Logging 中。

Dataproc documentation 中所述,要在 Logging 中启用作业驱动程序日志,请在创建集群时设置以下集群属性:

dataproc:dataproc.logging.stackdriver.job.driver.enable=true

以下集群属性也是必需的,并且在创建集群时默认设置:

dataproc:dataproc.logging.stackdriver.enable=true
dataproc:jobs.file-backed-output.enable=true

【讨论】:

  • 这对我来说感觉像是部分解决方案:我没有在集群内设置 Python 日志记录,但希望 Airflow 也能配置它。我希望有一些参数可以传递给 Airflow,告诉它设置 Python 日志记录,但这可能不是一个已实现的功能......
  • 我现在收到日志消息(在 GCP AI Platform 日志中,但不在 Airflow 日志中),使用 both print() 和 logger=logging.getLogger('projects/MyProjectName/logs/master-replica-0') 我没有认为我改变了任何东西,但我也很惊讶我之前没有看到我的测试消息...... sigh 要接受这个答案......
猜你喜欢
  • 2019-02-21
  • 2020-05-08
  • 2017-02-28
  • 1970-01-01
  • 2021-04-17
  • 2019-05-18
  • 1970-01-01
  • 2022-12-15
  • 1970-01-01
相关资源
最近更新 更多