Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
34 commits
Select commit Hold shift + click to select a range
aefb584
Ported to gradle build
Jun 28, 2013
38c3c3d
Merge pull request #24 from divchenko/feature/gradle-build
jhartman Jul 5, 2013
25f4460
FIX: change version back to 0.6.33
Jul 11, 2013
b718e17
Merge pull request #26 from vikstrous/0.6.33
jhartman Jul 11, 2013
11d6c89
Adjusting the norbert tuning parameters affected by the change from m…
Jul 15, 2013
ad21559
Merge remote-tracking branch 'upstream/master'
Jul 15, 2013
37ce05b
Upping the version
Jul 15, 2013
079ebea
Merge pull request #27 from abhinigam/master
jhartman Jul 16, 2013
1977d54
Updating Netty Version in Norbert for USCP 2050
Aug 1, 2013
38bcb5b
Merge pull request #1 from navina/master
jhartman Aug 5, 2013
276fb1f
Upping the version
Aug 5, 2013
92e1a21
Merge pull request #2 from navina/master
navina Aug 5, 2013
44d86ea
Modifying the JMX stats so that they remain unchanged even after mill…
Aug 12, 2013
135ebc2
Adding metrics: active pool size, current pool size, queue time, tota…
Aug 14, 2013
d981f34
Testing push mechanism
Aug 15, 2013
b43b725
Moving queueTime from endRequest to beginRequest, leaving existing la…
Sep 9, 2013
ef742cc
Fixing tests
Sep 9, 2013
f06acf7
Upgrading the version of norbert
Sep 9, 2013
0bd3b00
Added in support for selective retry support
Oct 22, 2013
74d5009
Making fixes based on code review comments
Oct 24, 2013
47ec2b9
Merge branch 'master' of https://github.com/linkedin/norbert into Fix…
Oct 25, 2013
25da9c3
Fixing the race condition which occurs if a node is started quickly a…
Nov 9, 2013
b868e39
Fixing tab ordering
Nov 9, 2013
f5578eb
Expose duplicatesOk configuration and add support for ListenerFuture
Nov 19, 2013
8cb3b2f
Remove reliance on a threadpool for callback execution
Nov 20, 2013
363b723
Chaning the name of the interface to PromiseListener
Nov 20, 2013
9ccbd87
Improving the exception handling to be more consistent
Nov 20, 2013
465bbf7
Cleaning up comments
Nov 21, 2013
7b48470
Hooking in FutureAdapterListener in all the places where FutureAdapte…
Dec 2, 2013
5294f84
Adding additional logging, handling expected zookeeper events
Dec 3, 2013
75e49db
Merge branch 'FixBugs' of https://github.com/linkedin/norbert into Fi…
Dec 3, 2013
8fa493f
Adding support for customizing behavior of overriding duplicatesOk be…
Dec 5, 2013
ba7229f
Merge branch 'FixBugs' of https://github.com/linkedin/norbert into Fi…
Dec 5, 2013
e8825c4
Creating a new API instead of modifying an existing API
Dec 7, 2013
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -3,3 +3,9 @@ lib_managed/
src_managed/
project/boot/
project/plugins/project
out/
build/
.gradle
*.iml
*.ipr
*.iws
1 change: 1 addition & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -269,3 +269,4 @@ If you are building a partitioned cluster then you will want to use the `Partiti
* maxConnectionsPerNode - the maximum number of open connections to a node. The total number of connections that can be opened by a network client is maxConnectionsPerNode * number of nodes
* staleRequestTimeoutMins - the number of minutes to keep a request that is waiting for a response
* staleRequestCleanupFrequenceMins - the frequency to clean up stale requests

27 changes: 27 additions & 0 deletions build.gradle
Original file line number Diff line number Diff line change
@@ -0,0 +1,27 @@
project.ext.isDefaultEnvironment = !project.hasProperty('overrideBuildEnvironment')

File getEnvironmentScript()
{
final File env = file(isDefaultEnvironment ? 'defaultEnvironment.gradle' : project.overrideBuildEnvironment)
assert env.isFile() : "The environment script [$env] does not exists or is not a file."
return env
}

apply from: environmentScript

project.ext.externalDependency = [
'zookeeper':'org.apache.zookeeper:zookeeper:3.3.0',
'protobuf':'com.google.protobuf:protobuf-java:2.4.0a',
'log4j':'log4j:log4j:1.2.16',
'netty':'io.netty:netty:3.5.11.Final',
'slf4jApi':'org.slf4j:slf4j-api:1.5.6',
'slf4jLog4j':'org.slf4j:slf4j-log4j12:1.5.6',
'specs':'org.scala-tools.testing:specs_2.8.1:1.6.8',
'mockitoAll':'org.mockito:mockito-all:1.8.4',
'cglib':'cglib:cglib:2.1_3',
'objenesis':'org.objenesis:objenesis:1.2',
'scalaCompiler': 'org.scala-lang:scala-compiler:2.8.1',
'scalaLibrary': 'org.scala-lang:scala-library:2.8.1',
'scalatest': 'org.scalatest:scalatest:1.2'
];

18 changes: 18 additions & 0 deletions cluster/build.gradle
Original file line number Diff line number Diff line change
@@ -0,0 +1,18 @@
apply plugin: 'java'
apply plugin: 'scala'

dependencies {
compile externalDependency.scalaLibrary
compile externalDependency.zookeeper
compile externalDependency.protobuf
compile externalDependency.log4j

testCompile externalDependency.specs
testCompile externalDependency.mockitoAll
testCompile externalDependency.cglib
testCompile externalDependency.objenesis

scalaTools externalDependency.scalaCompiler
scalaTools externalDependency.scalaLibrary
}

Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,8 @@ trait ZooKeeperClusterManagerComponent extends ClusterManagerComponent {
//synchronization is not needed between this method
case class NodeChildrenChanged(path: String) extends ZooKeeperMessage
case class NodeDataChanged(path: String) extends ZooKeeperMessage
case class NodeDeleted(path: String) extends ZooKeeperMessage
case class NodeCreated(path: String) extends ZooKeeperMessage
}

class ZooKeeperClusterManager(connectString: String, sessionTimeout: Int, serviceName: String)
Expand Down Expand Up @@ -90,6 +92,19 @@ trait ZooKeeperClusterManagerComponent extends ClusterManagerComponent {
handleCapabilityMemberChanged(path)
}

case NodeDeleted(path) => if (path.equals(MEMBERSHIP_NODE) || path.equals(AVAILABILITY_NODE)) {
//zookeeper data corrupted
log.fatal("Received a zookeeper corruption message")
} else if (path.startsWith(MEMBERSHIP_NODE) || path.startsWith(AVAILABILITY_NODE)) {
log.info("Received an event where a node was deleted:%s".format(path))
} else {
log.error("Node deleted unexpectedly %s".format(path))
}

case NodeCreated(path) => {
log.error("Received an unexpected create event:%s".format(path))
}

case m => log.error("Received unknown message: %s".format(m))
}
}
Expand Down Expand Up @@ -461,6 +476,10 @@ trait ZooKeeperClusterManagerComponent extends ClusterManagerComponent {
case EventType.NodeChildrenChanged => zooKeeperManager ! NodeChildrenChanged(event.getPath)

case EventType.NodeDataChanged => zooKeeperManager ! NodeDataChanged(event.getPath)

case EventType.NodeCreated => zooKeeperManager ! NodeCreated(event.getPath)

case EventType.NodeDeleted => zooKeeperManager ! NodeDeleted(event.getPath)
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -71,23 +71,131 @@ class FinishedRequestTimeTracker(clock: Clock, interval: Long) {
}
}

class TotalRequestProcessingTime[KeyT](clock:Clock, interval:Long) {
private val q = new java.util.concurrent.ConcurrentLinkedQueue[(Long, Long)]()
private val currentlyCleaning = new java.util.concurrent.atomic.AtomicBoolean

private def clean {
// Let only one thread clean at a time
if(currentlyCleaning.compareAndSet(false, true)) {
clean0
currentlyCleaning.set(false)
}
}

private def clean0 {
while(!q.isEmpty) {
val head = q.peek
if(head == null)
return

val (completion, processingTime) = head
if(clock.getCurrentTimeOffsetMicroseconds - completion > interval) {
q.remove(head)
} else {
return
}
}
}

def addTime(processingTime: Long) {
clean
q.offer( (clock.getCurrentTimeOffsetMicroseconds, processingTime) )
}

def getArray: Array[(Long, Long)] = {
clean
q.toArray(Array.empty[(Long, Long)])
}

def getTimings: Array[Long] = {
getArray.map(_._2).sorted
}

def total = {
getTimings.sum
}

def reset {
q.clear
}
}

class QueueTimeTracker[KeyT](clock: Clock, interval: Long) {
private val q = new java.util.concurrent.ConcurrentLinkedQueue[(Long, Long)]()
private val currentlyCleaning = new java.util.concurrent.atomic.AtomicBoolean

private def clean {
// Let only one thread clean at a time
if(currentlyCleaning.compareAndSet(false, true)) {
clean0
currentlyCleaning.set(false)
}
}

private def clean0 {
while(!q.isEmpty) {
val head = q.peek
if(head == null)
return

val (completion, processingTime) = head
if(clock.getCurrentTimeOffsetMicroseconds - completion > interval) {
q.remove(head)
} else {
return
}
}
}

def addTime(processingTime: Long) {
clean
q.offer( (clock.getCurrentTimeOffsetMicroseconds, processingTime) )
}

def getArray: Array[(Long, Long)] = {
clean
q.toArray(Array.empty[(Long, Long)])
}

def getTimings: Array[Long] = {
getArray.map(_._2).sorted
}

def total = {
getTimings.sum
}

def reset {
q.clear
}
}

// Threadsafe
class PendingRequestTimeTracker[KeyT](clock: Clock) {
private val numRequests = new AtomicInteger()

private val map : java.util.concurrent.ConcurrentMap[KeyT, Long] =
new java.util.concurrent.ConcurrentHashMap[KeyT, Long]

private val mapQueueTime : java.util.concurrent.ConcurrentMap[KeyT, Long] =
new java.util.concurrent.ConcurrentHashMap[KeyT, Long]

def getStartTime(key: KeyT) = Option(map.get(key))

def beginRequest(key: KeyT) {
//pre-condition for this method is the above method returns some
def getQueueTime(key: KeyT) = map.get(key)

def beginRequest(key: KeyT, queueTime: Long) {
numRequests.incrementAndGet
val now = clock.getCurrentTimeOffsetMicroseconds
map.put(key, now)
mapQueueTime.put(key, queueTime)
}

def endRequest(key: KeyT) {
map.remove(key)
mapQueueTime.remove(key)
}

def getTimings = {
Expand All @@ -108,20 +216,29 @@ class PendingRequestTimeTracker[KeyT](clock: Clock) {
class RequestTimeTracker[KeyT](clock: Clock, interval: Long) {
val finishedRequestTimeTracker = new FinishedRequestTimeTracker(clock, interval)
val pendingRequestTimeTracker = new PendingRequestTimeTracker[KeyT](clock)
val queueTimeTracker = new QueueTimeTracker[KeyT](clock, interval)//TODO
val totalRequestProcessingTimeTracker = new TotalRequestProcessingTime[KeyT](clock, interval)

def beginRequest(key: KeyT) {
pendingRequestTimeTracker.beginRequest(key)
def beginRequest(key: KeyT, queueTime: Long = 0) {
pendingRequestTimeTracker.beginRequest(key, queueTime)
}

def endRequest(key: KeyT) {
pendingRequestTimeTracker.getStartTime(key).foreach { startTime =>
//over time we will retire this since this does not account for the amount of time the request
//was stuck in the queue
finishedRequestTimeTracker.addTime(clock.getCurrentTimeOffsetMicroseconds - startTime)
val queueTime = pendingRequestTimeTracker.getQueueTime(key)
queueTimeTracker.addTime(queueTime)
totalRequestProcessingTimeTracker.addTime(queueTime + clock.getCurrentTimeOffsetMicroseconds - startTime)
}
pendingRequestTimeTracker.endRequest(key)
}

def reset {
finishedRequestTimeTracker.reset
pendingRequestTimeTracker.reset
queueTimeTracker.reset
totalRequestProcessingTimeTracker.reset
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -40,10 +40,10 @@ class AverageTimeTrackerSpec extends Specification {

(0 until 10).foreach { i =>
MockClock.currentTime = 1000L * i
tracker.beginRequest(i)
tracker.beginRequest(i,0)
(tracker.total / (i + 1)) must be_==(1000L * i / 2)
}
}

}
}
}
6 changes: 6 additions & 0 deletions defaultEnvironment.gradle
Original file line number Diff line number Diff line change
@@ -0,0 +1,6 @@
subprojects {
repositories {
mavenCentral()
}
}

14 changes: 7 additions & 7 deletions examples/src/main/resources/log4j.properties
Original file line number Diff line number Diff line change
Expand Up @@ -3,11 +3,11 @@ log4j.appender.CONSOLE=org.apache.log4j.ConsoleAppender
log4j.appender.CONSOLE.layout=org.apache.log4j.PatternLayout
log4j.appender.CONSOLE.layout.ConversionPattern=%d{ISO8601} - %-5p [%t:%C{1}@%L] - %m%n

#log4j.appender.R=org.apache.log4j.RollingFileAppender
#log4j.appender.R.File=/tmp/norbert.log
#log4j.appender.R.MaxFileSize=100KB
#log4j.appender.R.MaxBackupIndex=1
#log4j.appender.R.layout=org.apache.log4j.PatternLayout
#log4j.appender.R.layout.ConversionPattern=%p %t %c - %m%n
log4j.appender.R=org.apache.log4j.RollingFileAppender
log4j.appender.R.File=/tmp/norbert.log
log4j.appender.R.MaxFileSize=100KB
log4j.appender.R.MaxBackupIndex=1
log4j.appender.R.layout=org.apache.log4j.PatternLayout
log4j.appender.R.layout.ConversionPattern=%p %t %c - %m%n

#log4j.category.com.linkedin.norbert=debug
log4j.category.com.linkedin.norbert=debug
1 change: 1 addition & 0 deletions gradle.properties
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
version=0.6.36
11 changes: 11 additions & 0 deletions java-cluster/build.gradle
Original file line number Diff line number Diff line change
@@ -0,0 +1,11 @@
apply plugin: 'java'
apply plugin: 'scala'

dependencies {
compile project(':cluster')
compile externalDependency.scalaLibrary

scalaTools externalDependency.scalaCompiler
scalaTools externalDependency.scalaLibrary
}

Original file line number Diff line number Diff line change
Expand Up @@ -28,16 +28,16 @@ package object javacompat {
implicit def javaSetToImmutableSet[T](nodes: java.util.Set[T]): Set[T] = {
collection.JavaConversions.asScalaSet(nodes).foldLeft(Set[T]()) { (set, n) => set + n }
}

implicit def javaIntegerSetToScalaIntSet(set: java.util.Set[java.lang.Integer]): Set[Int] = {
collection.JavaConversions.asScalaSet(set).foldLeft(collection.immutable.Set.empty[Int]) { _ + _.intValue }
}

implicit def scalaIntSetToJavaIntegerSet(set: Set[Int]): java.util.Set[java.lang.Integer] = {
val result = new java.util.HashSet[java.lang.Integer](set.size)
set.foreach (result add _)
result
}
}

implicit def scalaNodeToJavaNode(node: SNode): JNode = {
if (node == null) null else JavaNode(node)
Expand Down
12 changes: 12 additions & 0 deletions java-network/build.gradle
Original file line number Diff line number Diff line change
@@ -0,0 +1,12 @@
apply plugin: 'java'
apply plugin: 'scala'

dependencies {
compile project(':network')
compile project(':java-cluster')
compile externalDependency.scalaLibrary

scalaTools externalDependency.scalaCompiler
scalaTools externalDependency.scalaLibrary
}

Loading