【问题标题】:Big query cron job without webservice没有 Web 服务的大查询 cron 作业
【发布时间】:2018-03-26 02:31:32
【问题描述】:

是否可以在不运行谷歌应用引擎网络服务的情况下运行处理用户数据的脚本?

使用较小的脚本效果很好,但是当我的脚本持续大约 40 分钟时,我收到错误:DeadlineExceededError

我的临时解决方法是在 windows VM 上使用 windows 调度程序并在命令行中使用 python 脚本

编辑:添加代码

jobs = []
jobs_status = []
jobs_error = []
# The project id whose datasets you'd like to list
PROJECT_NUMBER = 'project'
scope = ('https://www.googleapis.com/auth/bigquery',
         'https://www.googleapis.com/auth/cloud-platform',
         'https://www.googleapis.com/auth/drive',
         'https://spreadsheets.google.com/feeds')

credentials = ServiceAccountCredentials.from_json_keyfile_name('client_secrets.json', scope)

# Create the bigquery api client
service = googleapiclient.discovery.build('bigquery', 'v2', credentials=credentials)

def load_logs(source):
    body = {"rows": [
        {"json": source}
    ]}

    response = service.tabledata().insertAll(
        projectId=PROJECT_NUMBER,
        datasetId='test',
        tableId='test_log',
        body=body).execute()
    return response

def job_status():
    for job in jobs:
        _jobId = job['jobReference']['jobId']
        status = service.jobs().get(projectId=PROJECT_NUMBER, jobId=_jobId).execute()
        jobs_status.append(status['status']['state'])
        if 'errors' in status['status'].keys():
            query = str(status['configuration']['query']['query'])
            message = str(status['status']['errorResult']['message'])
            jobs_error.append({"query": query, "message": message})
    return jobs_status


def check_statues():
    while True:
        if all('DONE' in job for job in job_status()):
            return


def insert(query, tableid, disposition):
    job_body = {
     "configuration": {
      "query": {
       "query": query,
       "useLegacySql": True,
       "destinationTable": {
        "datasetId": "test",
        "projectId": "project",
        "tableId": tableid
       },
       "writeDisposition": disposition
      }
     }
    }

    r = service.jobs().insert(
        projectId=PROJECT_NUMBER,
        body=job_body).execute()
    jobs.append(r)
    return r



class MainPage(webapp2.RequestHandler):
    def get(self):
        query = "SELECT * FROM [gdocs_users.user_empty]"
        insert(query, 'users_data_p1', "WRITE_TRUNCATE")
        check_statues()
        query = "SELECT * FROM [gdocs_users.user_empty]"
        insert(query, 'users_data_p2', "WRITE_TRUNCATE")
        query = "SELECT * FROM [gdocs_users.user_%s]"
        for i in range(1, 1000):
            if i <= 600:
                insert(query % str(i).zfill(4), 'users_data_p1', "WRITE_APPEND")
            else:
                insert(query % str(i).zfill(4), 'user_data_p2', "WRITE_APPEND")
        for error in jobs_error:
            load_logs(error)


app = webapp2.WSGIApplication([
    ('/', MainPage),
], debug=True)

【问题讨论】:

    标签: python google-app-engine cron google-bigquery


    【解决方案1】:

    大多数 BigQuery 操作可以异步运行。你能告诉我们你的代码吗?

    例如,来自 Python BigQuery 文档:

    def query(query):
        client = bigquery.Client()
        query_job = client.run_async_query(str(uuid.uuid4()), query)
    
        query_job.begin()
        query_job.result()  # Wait for job to complete
    

    这是一个异步作业,代码选择等待查询完成。无需等待,而是在begin() 之后获取作业 ID。您可以使用Task Queue 将任务排入队列以便稍后运行,以检查该作业的结果。

    【讨论】:

    • 我必须创建服务,因为我有联合表(来自 gsuit),这些表不能使用 python API,正如你所见,我通过插入作业异步运行它们
    • 是的 - 查看您不应该使用的代码 check_statues()。然后,该函数将在插入作业后立即返回。然后稍后再回来检查该作业 ID 的结果。
    • 但是我在下面有一些查询,当我的插入完成时我正在运行它们。正如我在文档中读到的,在 60 秒后请求超时后会引发错误,所以问题是应用引擎是此类脚本的好地方,还是更适合小型后端的东西?
    • 是的。应用引擎很好。只需在异步脚本中运行即可。
    • 循环后我在脚本中使用了几次check_statues。再次处理有错误的查询,将它们保存到带有日志的表中。然后我运行几个不同的查询,这些查询需要完成以前的工作。我将缩放更改为基本并且问题消失了,因为我的脚本持续了大约 40 分钟。
    【解决方案2】:

    默认情况下,App Engine 服务使用automatic scaling,它对 HTTP 请求有 60 秒的限制,对任务队列请求有 10 分钟的限制。如果您将服务更改为使用基本或手动扩展,那么您的任务队列请求最多可以运行 24 小时。

    听起来您可能只需要一个实例来完成这项工作,因此除了默认服务之外,也许还需要创建第二个 service。在子文件夹中创建一个 bqservice 文件夹,其中包含以下 app.yaml 设置,这些设置使用最多一个实例的基本缩放:

    # bqsservice/app.yaml
    # Possibly use a separate service for your BQ code than
    # the rest of your app:
    service: bqservice
    runtime: python27
    api_version: 1
    # Keep low memory/cost B1 class?
    instance_class: B1
    # Limit max services to 1 to keep costs down. There is an
    # 8 instance hour limit to the free tier. This option still
    # scales to 0 when not in use.
    basic_scaling:
      max_instances: 1
    
    # Handlers:
    handlers:
    - url: /.*
      script: main.app
    

    然后在同一个服务中创建一个cron.yaml 来安排你的脚本运行。使用上面的示例配置,您可以将 BigQuery 逻辑放入一个 main.py 文件中,并在其中定义一个 WSGI 应用程序:

    # bqservice/main.py
    import webapp2
    
    class CronHandler(webapp2.RequestHandler):
    
        def post(self):
           # Handle your cron work
           # ....
    
    app = webapp2.WSGIApplication([
        #('/', MainPage),  # If you needed other handlers
        ('/mycron', CronHandler),
    ], debug=True)
    

    如果您不打算将 App Engine 应用用于其他任何用途,您可以将这一切都用于默认服务。如果您在默认服务之外执行此操作,则需要先将某些内容部署到默认服务,即使它只是带有静态文件的简单 app.yaml

    【讨论】:

    • 谢谢布雷特J!此解决方案可能有效 - 但请注意 App Engine 在等待 BigQuery 返回时没有执行任何操作。更好的方法是在作业入队后立即返回,然后稍后再回来查看作业是否完成。
    • @BrettJ 谢谢你的帮助。我还想添加如果有人会尝试,您必须在您的 cron.yaml 中添加 target: your_service 以确定您要使用的服务
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2014-03-27
    • 1970-01-01
    • 2014-05-26
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多