-
Notifications
You must be signed in to change notification settings - Fork 8
/
Sender.java
81 lines (66 loc) · 2.49 KB
/
Sender.java
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
package com.kickstarter.dropwizard.metrics.influxdb.io;
import com.google.common.annotations.VisibleForTesting;
import com.google.common.collect.EvictingQueue;
import com.kickstarter.dropwizard.metrics.influxdb.InfluxDbMeasurement;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.io.UnsupportedEncodingException;
import java.util.Collection;
import static java.util.stream.Collectors.toList;
/**
* Sends measurements to InfluxDB. Uses an {@link EvictingQueue} to store and retry measurements that have
* failed to send, and timestamps measurements at the configured {@code precision}, up to millisecond precision.
*/
public class Sender {
private static final Logger log = LoggerFactory.getLogger(InfluxDbTcpWriter.class);
public static final int DEFAULT_QUEUE_SIZE = 5000;
private static final String SEPARATOR = "\n";
private final InfluxDbWriter writer;
private final EvictingQueue<InfluxDbMeasurement> queuedInfluxDbMeasurements;
public Sender(final InfluxDbWriter writer) {
this(writer, DEFAULT_QUEUE_SIZE);
}
public Sender(final InfluxDbWriter writer, final int queueSize) {
this.writer = writer;
this.queuedInfluxDbMeasurements = EvictingQueue.create(queueSize);
}
@VisibleForTesting int queuedMeasures() {
return queuedInfluxDbMeasurements.size();
}
/**
* Sends the provided {@link InfluxDbMeasurement measurements} to InfluxDB.
*
* @return true if the measurements were successfully sent.
*/
public boolean send(final Collection<InfluxDbMeasurement> influxDbMeasurements) {
queuedInfluxDbMeasurements.addAll(influxDbMeasurements);
if (queuedInfluxDbMeasurements.isEmpty()) {
return true;
}
final String measureLines = String.join(
SEPARATOR,
queuedInfluxDbMeasurements.stream()
.map(InfluxDbMeasurement::toLine)
.collect(toList())
) + SEPARATOR;
try {
final byte[] measureBytes = measureLines.getBytes("UTF-8");
writer.writeBytes(measureBytes);
queuedInfluxDbMeasurements.clear();
return true;
} catch (final UnsupportedEncodingException e) {
log.warn("failed to send metrics", e);
} catch (final Exception e) {
log.warn("failed to send metrics", e);
try {
writer.close();
} catch (final Exception e2) {
log.warn("failed to close metrics connection", e2);
}
}
if (queuedInfluxDbMeasurements.remainingCapacity() == 0) {
log.warn("Queued measurements at capacity");
}
return false;
}
}