【问题标题】:How to write properly from an Apache Beam job to InfluxDB?如何从 Apache Beam 作业正确写入 InfluxDB?
【发布时间】:2019-09-10 06:08:08
【问题描述】:

我在 Google Cloud Dataflow 中执行了一个 Apache Beam (v2.15.0) 作业。该作业是一项流式作业,它从 Google Cloud PubSub 读取消息、对数据进行分组并尝试将它们写入 InfluxDB。

主要问题是当数据必须写入 InfluxDB 时,同时向 InfluxDB 发送了太多写入请求,我收到很多以下错误:

{
 insertId: "750692785846999589:29668:0:536239"  
 jsonPayload: {
  job: "2019-09-09_07_41_32-8614892132258718725"   
  logger: "functions.DashboardFunctions$SendPostLiteFn2"   
  message: "The writing to InfluxDB has failed. LOG: engine: error syncing wal"   
  stage: "P2"   
  step: "L-I-SendPost"   
  thread: "339"   
  work: "5a8|130|1|091443-992359958146824871"   
  worker: "dataflow-dashboard-test-09090741-ixek-harness-r7cz"   
 }
 labels: {
  compute.googleapis.com/resource_id: "750692785846999589"   
  compute.googleapis.com/resource_name: "dataflow-dashboard-test-09090741-ixek-harness-r7cz"   
  compute.googleapis.com/resource_type: "instance"   
  dataflow.googleapis.com/job_id: "2019-09-09_07_41_32-8614892132258718725"   
  dataflow.googleapis.com/job_name: "dataflow-dashboard-test"   
  dataflow.googleapis.com/region: "europe-west1"   
 }
 logName: "projects/XXXX/logs/dataflow.googleapis.com%2Fworker"  
 receiveTimestamp: "2019-09-09T14:46:50.419635294Z"  
 resource: {
  labels: {
   job_id: "2019-09-09_07_41_32-8614892132258718725"    
   job_name: "dataflow-dashboard-test"    
   project_id: "XXXX"    
   region: "europe-west1"    
   step_id: "L-I-SendPost"    
  }
  type: "dataflow_step"   
 }
 severity: "ERROR"  
 timestamp: "2019-09-09T14:46:35.640Z"  
}

我不理解这种行为,因为我认为我正在批量写入...

您可以在下面找到我编写的将数据批量发送到 InfluxDB 的代码:

/**
 * A DoFn that sends POSTS to InfluxDB.
 */
public static class WrtieToInfluxDB extends DoFn<List<String>, Void> {
    private static final Logger LOG = LoggerFactory.getLogger(WrtieToInfluxDB.class);

    private InfluxDB influxClient = null;
    private BatchPoints batchPoints = null;

    @Setup
    public void initInfluxClient() {
        OkHttpClient.Builder okHttpClientBuilder = new OkHttpClient().newBuilder()
                .connectTimeout(40, TimeUnit.SECONDS)
                .readTimeout(60, TimeUnit.SECONDS)
                .writeTimeout(60, TimeUnit.SECONDS);

        this.influxClient =
                InfluxDBFactory
                        .connect(
                                Settings.getInfluxUrl(),
                                Settings.getInfluxUsername(),
                                Settings.getInfluxPassword(),
                                okHttpClientBuilder)
                        .enableGzip();
    }

    @StartBundle
    public void startBundle(StartBundleContext context) {
        Util.MyOptions myOptions = context.getPipelineOptions().as(Util.MyOptions.class);

        this.batchPoints =
                BatchPoints
                        .database(Settings.getDbNameProvider(Settings.StatsType.DASHBOARD, myOptions.getTest()))
                        .consistency(InfluxDB.ConsistencyLevel.ONE)
                        .build();
    }

    @ProcessElement
    public void processElement(final ProcessContext c) {
        // Code to builds the point
        Point.Builder point = Point
            ...

        // Inserts the point into the batch object
        this.batchPoints.point(point.build());
    }

    @FinishBundle
    public void finishBundle(FinishBundleContext context) {
        // Send a POST to influxDB with all the points of the batchPoints
        try {
            this.influxClient.write(this.batchPoints);
            this.influxClient.flush();
        } catch (Throwable t) {
            LOG.error("The writing to InfluxDB has failed. LOG: " + t.getMessage());
        }
    }
}

编辑:作业有一个 1 分钟的固定窗口。

【问题讨论】:

    标签: java google-cloud-dataflow apache-beam influxdb


    【解决方案1】:

    数据流流尝试通过使用较小的包大小来并行化您的工作,这意味着可以在非常小的包(甚至小到单个元素/KV 对)上调用@FinishBundle。

    您应该查看 GroupIntoBatches 以帮助您增加批处理大小,这应该会减少您在 InfluxDB 上放置的负载量。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2020-05-01
      • 1970-01-01
      • 2020-12-21
      • 2021-10-10
      • 2017-09-03
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多