mirror of https://github.com/apache/druid.git
parent
ff56573910
commit
faa74ce2cb
|
@ -313,10 +313,12 @@ public class RealtimePlumber implements Plumber
|
||||||
final Interval interval = sink.getInterval();
|
final Interval interval = sink.getInterval();
|
||||||
|
|
||||||
for (FireHydrant hydrant : sink) {
|
for (FireHydrant hydrant : sink) {
|
||||||
if (!hydrant.hasSwapped()) {
|
synchronized (hydrant) {
|
||||||
log.info("Hydrant[%s] hasn't swapped yet, swapping. Sink[%s]", hydrant, sink);
|
if (!hydrant.hasSwapped()) {
|
||||||
final int rowCount = persistHydrant(hydrant, schema, interval);
|
log.info("Hydrant[%s] hasn't swapped yet, swapping. Sink[%s]", hydrant, sink);
|
||||||
metrics.incrementRowOutputCount(rowCount);
|
final int rowCount = persistHydrant(hydrant, schema, interval);
|
||||||
|
metrics.incrementRowOutputCount(rowCount);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
Loading…
Reference in New Issue