【问题标题】:Dataflow stops streaming to BigQuery without errorsDataflow 停止流式传输到 BigQuery 且没有错误
【发布时间】:2018-12-04 10:30:44
【问题描述】:

我们开始使用 Dataflow 从 PubSub 和 Stream 读取到 BigQuery。 数据流应该 24/7 全天候工作,因为 pubsub 会不断更新世界各地多个网站的分析数据。

代码如下:

from __future__ import absolute_import

import argparse
import json
import logging

import apache_beam as beam
from apache_beam.io import ReadFromPubSub, WriteToBigQuery
from apache_beam.options.pipeline_options import PipelineOptions
from apache_beam.options.pipeline_options import SetupOptions

logger = logging.getLogger()

TABLE_IDS = {
    'table_1': 0,
    'table_2': 1,
    'table_3': 2,
    'table_4': 3,
    'table_5': 4,
    'table_6': 5,
    'table_7': 6,
    'table_8': 7,
    'table_9': 8,
    'table_10': 9,
    'table_11': 10,
    'table_12': 11,
    'table_13': 12
 }


def separate_by_table(element, num):
    return TABLE_IDS[element.get('meta_type')]


class ExtractingDoFn(beam.DoFn):
    def process(self, element):
        yield json.loads(element)


def run(argv=None):
    """Main entry point; defines and runs the wordcount pipeline."""
    logger.info('STARTED!')
    parser = argparse.ArgumentParser()
    parser.add_argument('--topic',
                        dest='topic',
                        default='projects/PROJECT_NAME/topics/TOPICNAME',
                        help='Gloud topic in form "projects/<project>/topics/<topic>"')
    parser.add_argument('--table',
                        dest='table',
                        default='PROJECTNAME:DATASET_NAME.event_%s',
                        help='Gloud topic in form "PROJECT:DATASET.TABLE"')
    known_args, pipeline_args = parser.parse_known_args(argv)

    # We use the save_main_session option because one or more DoFn's in this
    # workflow rely on global context (e.g., a module imported at module level).
    pipeline_options = PipelineOptions(pipeline_args)
    pipeline_options.view_as(SetupOptions).save_main_session = True
    p = beam.Pipeline(options=pipeline_options)

    lines = p | ReadFromPubSub(known_args.topic)
    datas = lines | beam.ParDo(ExtractingDoFn())
    by_table = datas | beam.Partition(separate_by_table, 13)

    # Create a stream for each table
    for table, id in TABLE_IDS.items():
        by_table[id] | 'write to %s' % table >> WriteToBigQuery(known_args.table % table)

    result = p.run()
    result.wait_until_finish()


if __name__ == '__main__':
    logger.setLevel(logging.INFO)
    run()

它工作正常,但一段时间后(2-3 天)它会因某种原因停止流式传输。 当我检查作业状态时,它在日志部分中不包含任何错误(您知道,在数据流的作业详细信息中标有红色“!”)。如果我取消作业并再次运行它 - 它会像往常一样再次开始工作。 如果我检查 Stackdriver 是否有其他日志,以下是发生的所有错误: 以下是作业执行时定期发生的一些警告: 其中之一的详细信息:

 {
 insertId: "397122810208336921:865794:0:479132535"  

jsonPayload: {
  exception: "java.lang.IllegalStateException: Cannot be called on unstarted operation.
    at com.google.cloud.dataflow.worker.fn.data.RemoteGrpcPortWriteOperation.getElementsSent(RemoteGrpcPortWriteOperation.java:111)
    at com.google.cloud.dataflow.worker.fn.control.BeamFnMapTaskExecutor$SingularProcessBundleProgressTracker.updateProgress(BeamFnMapTaskExecutor.java:293)
    at com.google.cloud.dataflow.worker.fn.control.BeamFnMapTaskExecutor$SingularProcessBundleProgressTracker.periodicProgressUpdate(BeamFnMapTaskExecutor.java:280)
    at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511)
    at java.util.concurrent.FutureTask.runAndReset(FutureTask.java:308)
    at java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.access$301(ScheduledThreadPoolExecutor.java:180)
    at java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.run(ScheduledThreadPoolExecutor.java:294)
    at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1142)
    at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:617)
    at java.lang.Thread.run(Thread.java:745)
"   
  job: "2018-11-30_10_35_19-13557985235326353911"   
  logger: "com.google.cloud.dataflow.worker.fn.control.BeamFnMapTaskExecutor"   
  message: "Progress updating failed 4 times. Following exception safely handled."   
  stage: "S0"   
  thread: "62"   
  work: "c-8756541438010208464"   
  worker: "beamapp-vitar-1130183512--11301035-mdna-harness-lft7"   
 }

labels: {
  compute.googleapis.com/resource_id: "397122810208336921"   
  compute.googleapis.com/resource_name: "beamapp-vitar-1130183512--11301035-mdna-harness-lft7"   
  compute.googleapis.com/resource_type: "instance"   
  dataflow.googleapis.com/job_id: "2018-11-30_10_35_19-13557985235326353911"   
  dataflow.googleapis.com/job_name: "beamapp-vitar-1130183512-742054"   
  dataflow.googleapis.com/region: "europe-west1"   
 }
 logName: "projects/PROJECTNAME/logs/dataflow.googleapis.com%2Fharness"  
 receiveTimestamp: "2018-12-03T20:33:00.444208704Z"  

resource: {

labels: {
   job_id: "2018-11-30_10_35_19-13557985235326353911"    
   job_name: "beamapp-vitar-1130183512-742054"    
   project_id: PROJECTNAME
   region: "europe-west1"    
   step_id: ""    
  }
  type: "dataflow_step"   
 }
 severity: "WARNING"  
 timestamp: "2018-12-03T20:32:59.442Z"  
}

这是它似乎开始出现问题的时刻: 可能有帮助的其他信息消息:

根据这些消息,我们没有耗尽内存/处理能力等。作业使用以下参数运行:

python -m start --streaming True --runner DataflowRunner --project PROJECTNAME --temp_location gs://BUCKETNAME/tmp/ --region europe-west1 --disk_size_gb 30 --machine_type n1-standard-1 --use_public_ips false --num_workers 1 --max_num_workers 1 --autoscaling_algorithm NONE

这可能是什么问题?

【问题讨论】:

    标签: google-cloud-platform google-cloud-dataflow apache-beam google-cloud-pubsub


    【解决方案1】:

    这并不是真正的答案,更有助于确定原因:到目前为止,我使用 python SDK 启动的所有流式 Dataflow 作业在几天后都以这种方式停止,无论它们是否使用 BigQuery 作为接收器。所以原因似乎是streaming jobs with the python SDK are still in beta.

    我的个人解决方案:使用 Dataflow 模板从 Pub/Sub 流式传输到 BigQuery(从而避免使用 Python SDK),然后在 BigQuery 中安排查询以定期处理数据。不幸的是,这可能不适合您的用例。

    【讨论】:

      【解决方案2】:

      在我的公司中,我们遇到了同样的问题,正如 OP 所描述的,具有类似的用例。

      不幸的是,这个问题是真实的、具体的,而且显然是随机发生的。

      作为一种解决方法,我们正在考虑使用 java SDK 重写我们的管道。

      【讨论】:

      • 您好,同样的问题。你用Java重写管道解决了吗?
      • 你好@lordcenzin,我们目前正在试验这个cloud.google.com/dataflow/docs/guides/templates/…。不过,我使用最新的 Python sdk (2.11.0) 执行了管道更新,它似乎表现良好。有问题的日志消失了。变更日志没有明确声明解决了这个问题,所以我的手指仍然交叉。
      • 在 2.15.0 上仍然可以看到这个
      【解决方案3】:

      我遇到了类似的问题,发现警告日志包含隐藏在提示错误的 java 日志中的 python Stack 跟踪。

      工作人员不断重试这些错误,导致它们崩溃并完全冻结管道。我最初认为工人数量太少,所以扩大工人数量,但管道冻结的时间更长。

      我在本地运行管道并将 pubsub 消息导出为文本,并确定它们包含脏数据(与 BQ 表架构不匹配的消息),并且由于我没有异常处理,这似乎是管道的原因冻结。

      添加函数仅接受第一个键与 BQ 架构的预期列匹配的记录解决了我的问题,并且数据流作业一直在运行,没有任何问题。

      def bad_records(row):
          if 'key1' in row:
              yield row
          else:
              print('bad row',row)
      
      
      |'exclude bad records' >> beam.ParDo(bad_records)
      

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 2018-06-21
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2015-07-15
        • 2021-12-02
        相关资源
        最近更新 更多