【发布时间】: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