Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
A set of AtomicInteger bound to the Core object that are incremented and decremented on buffer read and write allow us to indicate current utilization of the buffers. Using Metrics was not an option here, because the most precise calculation (one minute average) was way to deferred. The new CLI parameter -s/--statistics enables periodical printing of buffer utilization to STDOUT. Next step is to write this information to MongoDB and display it in the web interface. fixes #SERVER-199
- Loading branch information
Lennart Koopmann
committed
Nov 5, 2012
1 parent
9a498a2
commit c3bc81f
Showing
12 changed files
with
239 additions
and
8 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Original file line | Diff line number | Diff line change |
---|---|---|---|
@@ -0,0 +1,48 @@ | |||
/** | |||
* Copyright 2012 Lennart Koopmann <lennart@socketfeed.com> | |||
* | |||
* This file is part of Graylog2. | |||
* | |||
* Graylog2 is free software: you can redistribute it and/or modify | |||
* it under the terms of the GNU General Public License as published by | |||
* the Free Software Foundation, either version 3 of the License, or | |||
* (at your option) any later version. | |||
* | |||
* Graylog2 is distributed in the hope that it will be useful, | |||
* but WITHOUT ANY WARRANTY; without even the implied warranty of | |||
* MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the | |||
* GNU General Public License for more details. | |||
* | |||
* You should have received a copy of the GNU General Public License | |||
* along with Graylog2. If not, see <http://www.gnu.org/licenses/>. | |||
* | |||
*/ | |||
package org.graylog2.buffers; | |||
|
|||
import java.util.concurrent.atomic.AtomicInteger; | |||
import org.apache.log4j.Logger; | |||
|
|||
/** | |||
* @author Lennart Koopmann <lennart@socketfeed.com> | |||
*/ | |||
public class BufferWatermark { | |||
|
|||
private static final Logger LOG = Logger.getLogger(BufferWatermark.class); | |||
|
|||
private final int bufferSize; | |||
private final AtomicInteger watermark; | |||
|
|||
public BufferWatermark(int bufferSize, AtomicInteger watermark) { | |||
this.bufferSize = bufferSize; | |||
this.watermark = watermark; | |||
} | |||
|
|||
public int getUtilization() { | |||
return watermark.get(); | |||
} | |||
|
|||
public float getUtilizationPercentage() { | |||
return getUtilization()/bufferSize*100; | |||
} | |||
|
|||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
24 changes: 19 additions & 5 deletions
24
src/main/java/org/graylog2/initializers/AMQPSyncInitializer.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
47 changes: 47 additions & 0 deletions
47
src/main/java/org/graylog2/initializers/BufferWatermarkInitializer.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Original file line | Diff line number | Diff line change |
---|---|---|---|
@@ -0,0 +1,47 @@ | |||
/** | |||
* Copyright 2012 Lennart Koopmann <lennart@socketfeed.com> | |||
* | |||
* This file is part of Graylog2. | |||
* | |||
* Graylog2 is free software: you can redistribute it and/or modify | |||
* it under the terms of the GNU General Public License as published by | |||
* the Free Software Foundation, either version 3 of the License, or | |||
* (at your option) any later version. | |||
* | |||
* Graylog2 is distributed in the hope that it will be useful, | |||
* but WITHOUT ANY WARRANTY; without even the implied warranty of | |||
* MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the | |||
* GNU General Public License for more details. | |||
* | |||
* You should have received a copy of the GNU General Public License | |||
* along with Graylog2. If not, see <http://www.gnu.org/licenses/>. | |||
* | |||
*/ | |||
package org.graylog2.initializers; | |||
|
|||
import org.graylog2.Core; | |||
import org.graylog2.periodical.BufferWatermarkThread; | |||
|
|||
/** | |||
* @author Lennart Koopmann <lennart@socketfeed.com> | |||
*/ | |||
public class BufferWatermarkInitializer extends SimpleFixedRateScheduleInitializer implements Initializer { | |||
|
|||
public BufferWatermarkInitializer(Core graylogServer) { | |||
this.graylogServer = graylogServer; | |||
} | |||
|
|||
@Override | |||
public void initialize() { | |||
configureScheduler( | |||
new BufferWatermarkThread(this.graylogServer), | |||
BufferWatermarkThread.INITIAL_DELAY, | |||
BufferWatermarkThread.PERIOD | |||
); | |||
} | |||
|
|||
@Override | |||
public boolean masterOnly() { | |||
return true; | |||
} | |||
} |
73 changes: 73 additions & 0 deletions
73
src/main/java/org/graylog2/periodical/BufferWatermarkThread.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Original file line | Diff line number | Diff line change |
---|---|---|---|
@@ -0,0 +1,73 @@ | |||
/** | |||
* Copyright 2012 Lennart Koopmann <lennart@socketfeed.com> | |||
* | |||
* This file is part of Graylog2. | |||
* | |||
* Graylog2 is free software: you can redistribute it and/or modify | |||
* it under the terms of the GNU General Public License as published by | |||
* the Free Software Foundation, either version 3 of the License, or | |||
* (at your option) any later version. | |||
* | |||
* Graylog2 is distributed in the hope that it will be useful, | |||
* but WITHOUT ANY WARRANTY; without even the implied warranty of | |||
* MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the | |||
* GNU General Public License for more details. | |||
* | |||
* You should have received a copy of the GNU General Public License | |||
* along with Graylog2. If not, see <http://www.gnu.org/licenses/>. | |||
* | |||
*/ | |||
package org.graylog2.periodical; | |||
|
|||
import java.util.concurrent.atomic.AtomicInteger; | |||
import org.apache.log4j.Logger; | |||
import org.graylog2.Core; | |||
import org.graylog2.buffers.BufferWatermark; | |||
import org.joda.time.DateTime; | |||
|
|||
/** | |||
* @author Lennart Koopmann <lennart@socketfeed.com> | |||
*/ | |||
public class BufferWatermarkThread implements Runnable { | |||
|
|||
private static final Logger LOG = Logger.getLogger(BufferWatermarkThread.class); | |||
|
|||
public static final int INITIAL_DELAY = 0; | |||
public static final int PERIOD = 5; | |||
|
|||
private final Core graylogServer; | |||
|
|||
public BufferWatermarkThread(Core graylogServer) { | |||
this.graylogServer = graylogServer; | |||
} | |||
|
|||
@Override | |||
public void run() { | |||
checkValidity(graylogServer.processBufferWatermark()); | |||
checkValidity(graylogServer.outputBufferWatermark()); | |||
|
|||
int ringSize = graylogServer.getConfiguration().getRingSize(); | |||
|
|||
BufferWatermark oWm = new BufferWatermark(ringSize, graylogServer.outputBufferWatermark()); | |||
|
|||
BufferWatermark pWm = new BufferWatermark(ringSize, graylogServer.processBufferWatermark()); | |||
|
|||
if (graylogServer.isStatsMode()) { | |||
DateTime now = new DateTime(); | |||
System.out.println("[util] [" + now + "] OutputBuffer is at " | |||
+ oWm.getUtilizationPercentage() + "%. [" + oWm.getUtilization() + "/" + ringSize +"]"); | |||
System.out.println("[util] [" + now + "] ProcessBuffer is at " | |||
+ pWm.getUtilizationPercentage() + "%. [" + pWm.getUtilization() + "/" + ringSize +"]"); | |||
} | |||
} | |||
|
|||
private void checkValidity(AtomicInteger watermark) { | |||
// This should never happen, but just to make sure... | |||
int x = watermark.get(); | |||
if (x < 0) { | |||
LOG.warn("Reset a watermark to 0 because it was <" + x + ">"); | |||
watermark.set(0); | |||
} | |||
} | |||
|
|||
} |