Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
阿信在这里
SkyWalking
提交
fdaa3186
S
SkyWalking
项目概览
阿信在这里
/
SkyWalking
与 Fork 源项目一致
Fork自
山不在高_有仙则灵 / SkyWalking
通知
1
Star
0
Fork
0
代码
文件
提交
分支
Tags
贡献者
分支图
Diff
Issue
0
列表
看板
标记
里程碑
合并请求
0
Wiki
0
Wiki
分析
仓库
DevOps
项目成员
Pages
S
SkyWalking
项目概览
项目概览
详情
发布
仓库
仓库
文件
提交
分支
标签
贡献者
分支图
比较
Issue
0
Issue
0
列表
看板
标记
里程碑
合并请求
0
合并请求
0
Pages
分析
分析
仓库分析
DevOps
Wiki
0
Wiki
成员
成员
收起侧边栏
关闭侧边栏
动态
分支图
创建新Issue
提交
Issue看板
体验新版 GitCode,发现更多精彩内容 >>
提交
fdaa3186
编写于
12月 04, 2016
作者:
wu-sheng
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
add AbstractSpanDisruptor to avoid endless blocking of next()
上级
6fd3ce2d
变更
6
显示空白变更内容
内联
并排
Showing
6 changed file
with
67 addition
and
9 deletion
+67
-9
skywalking-storage-center/skywalking-routing/src/main/java/com/a/eye/skywalking/routing/disruptor/AbstractSpanDisruptor.java
...e/skywalking/routing/disruptor/AbstractSpanDisruptor.java
+39
-0
skywalking-storage-center/skywalking-routing/src/main/java/com/a/eye/skywalking/routing/disruptor/NoopSpanDisruptor.java
...a/eye/skywalking/routing/disruptor/NoopSpanDisruptor.java
+5
-0
skywalking-storage-center/skywalking-routing/src/main/java/com/a/eye/skywalking/routing/disruptor/SpanDisruptor.java
...com/a/eye/skywalking/routing/disruptor/SpanDisruptor.java
+1
-1
skywalking-storage-center/skywalking-routing/src/main/java/com/a/eye/skywalking/routing/disruptor/ack/AckSpanDisruptor.java
...ye/skywalking/routing/disruptor/ack/AckSpanDisruptor.java
+10
-3
skywalking-storage-center/skywalking-routing/src/main/java/com/a/eye/skywalking/routing/disruptor/request/RequestSpanDisruptor.java
...lking/routing/disruptor/request/RequestSpanDisruptor.java
+11
-3
skywalking-storage-center/skywalking-routing/src/main/java/com/a/eye/skywalking/routing/router/Router.java
...main/java/com/a/eye/skywalking/routing/router/Router.java
+1
-2
未找到文件。
skywalking-storage-center/skywalking-routing/src/main/java/com/a/eye/skywalking/routing/disruptor/AbstractSpanDisruptor.java
0 → 100644
浏览文件 @
fdaa3186
package
com.a.eye.skywalking.routing.disruptor
;
import
com.lmax.disruptor.InsufficientCapacityException
;
import
com.lmax.disruptor.RingBuffer
;
/**
* Created by wusheng on 2016/12/4.
*/
public
class
AbstractSpanDisruptor
{
private
volatile
boolean
isRunning
=
true
;
public
long
getRingBufferSequence
(
RingBuffer
buffer
){
long
sequence
;
while
(
true
)
{
try
{
if
(!
isRunning
){
return
-
1
;
}
sequence
=
buffer
.
tryNext
();
break
;
}
catch
(
InsufficientCapacityException
e
)
{
try
{
Thread
.
sleep
(
1L
);
}
catch
(
InterruptedException
e1
)
{
}
}
}
return
sequence
;
}
protected
boolean
isShutdown
(){
return
!
isRunning
;
}
protected
void
shutdown
(){
isRunning
=
false
;
}
}
skywalking-storage-center/skywalking-routing/src/main/java/com/a/eye/skywalking/routing/disruptor/NoopSpanDisruptor.java
浏览文件 @
fdaa3186
...
...
@@ -7,6 +7,11 @@ import com.a.eye.skywalking.network.grpc.RequestSpan;
* Created by xin on 2016/11/29.
*/
public
class
NoopSpanDisruptor
extends
SpanDisruptor
{
public
static
NoopSpanDisruptor
INSTANCE
=
new
NoopSpanDisruptor
();
private
NoopSpanDisruptor
()
{
}
@Override
public
boolean
saveSpan
(
AckSpan
ackSpan
)
{
...
...
skywalking-storage-center/skywalking-routing/src/main/java/com/a/eye/skywalking/routing/disruptor/SpanDisruptor.java
浏览文件 @
fdaa3186
...
...
@@ -48,7 +48,7 @@ public class SpanDisruptor {
}
public
void
shutdown
()
{
ack
SpanDisruptor
.
shutdown
();
request
SpanDisruptor
.
shutdown
();
ackSpanDisruptor
.
shutdown
();
}
}
skywalking-storage-center/skywalking-routing/src/main/java/com/a/eye/skywalking/routing/disruptor/ack/AckSpanDisruptor.java
浏览文件 @
fdaa3186
...
...
@@ -6,6 +6,8 @@ import com.a.eye.skywalking.logging.api.ILog;
import
com.a.eye.skywalking.logging.api.LogManager
;
import
com.a.eye.skywalking.network.grpc.AckSpan
;
import
com.a.eye.skywalking.routing.config.Config
;
import
com.a.eye.skywalking.routing.disruptor.AbstractSpanDisruptor
;
import
com.lmax.disruptor.InsufficientCapacityException
;
import
com.lmax.disruptor.RingBuffer
;
import
com.lmax.disruptor.dsl.Disruptor
;
import
com.lmax.disruptor.util.DaemonThreadFactory
;
...
...
@@ -13,7 +15,7 @@ import com.lmax.disruptor.util.DaemonThreadFactory;
/**
* Created by xin on 2016/11/29.
*/
public
class
AckSpanDisruptor
{
public
class
AckSpanDisruptor
extends
AbstractSpanDisruptor
{
private
static
ILog
logger
=
LogManager
.
getLogger
(
AckSpanDisruptor
.
class
);
private
Disruptor
<
AckSpanHolder
>
ackSpanDisruptor
;
private
RingBuffer
<
AckSpanHolder
>
ackSpanRingBuffer
;
...
...
@@ -29,7 +31,11 @@ public class AckSpanDisruptor {
}
public
boolean
saveAckSpan
(
AckSpan
ackSpan
)
{
long
sequence
=
ackSpanRingBuffer
.
next
();
long
sequence
=
this
.
getRingBufferSequence
(
ackSpanRingBuffer
);
if
(
this
.
isShutdown
()){
return
false
;
}
try
{
AckSpanHolder
data
=
ackSpanRingBuffer
.
get
(
sequence
);
data
.
setAckSpan
(
ackSpan
);
...
...
@@ -44,9 +50,10 @@ public class AckSpanDisruptor {
}
}
public
void
shutdown
()
{
super
.
shutdown
();
ackSpanEventHandler
.
stop
();
ackSpanDisruptor
.
shutdown
();
}
}
skywalking-storage-center/skywalking-routing/src/main/java/com/a/eye/skywalking/routing/disruptor/request/RequestSpanDisruptor.java
浏览文件 @
fdaa3186
...
...
@@ -6,6 +6,8 @@ import com.a.eye.skywalking.logging.api.ILog;
import
com.a.eye.skywalking.logging.api.LogManager
;
import
com.a.eye.skywalking.network.grpc.RequestSpan
;
import
com.a.eye.skywalking.routing.config.Config
;
import
com.a.eye.skywalking.routing.disruptor.AbstractSpanDisruptor
;
import
com.lmax.disruptor.InsufficientCapacityException
;
import
com.lmax.disruptor.RingBuffer
;
import
com.lmax.disruptor.dsl.Disruptor
;
import
com.lmax.disruptor.util.DaemonThreadFactory
;
...
...
@@ -13,7 +15,7 @@ import com.lmax.disruptor.util.DaemonThreadFactory;
/**
* Created by xin on 2016/11/29.
*/
public
class
RequestSpanDisruptor
{
public
class
RequestSpanDisruptor
extends
AbstractSpanDisruptor
{
private
static
ILog
logger
=
LogManager
.
getLogger
(
RequestSpanDisruptor
.
class
);
private
Disruptor
<
RequestSpanHolder
>
requestSpanDisruptor
;
private
RingBuffer
<
RequestSpanHolder
>
requestSpanRingBuffer
;
...
...
@@ -28,8 +30,13 @@ public class RequestSpanDisruptor {
}
public
boolean
saveRequestSpan
(
RequestSpan
requestSpan
)
{
long
sequence
=
requestSpanRingBuffer
.
next
();
long
sequence
=
this
.
getRingBufferSequence
(
requestSpanRingBuffer
);
if
(
this
.
isShutdown
()){
return
false
;
}
try
{
RequestSpanHolder
data
=
requestSpanRingBuffer
.
get
(
sequence
);
data
.
setRequestSpan
(
requestSpan
);
...
...
@@ -44,7 +51,8 @@ public class RequestSpanDisruptor {
}
}
public
void
shutDown
()
{
public
void
shutdown
()
{
super
.
shutdown
();
eventHandler
.
stop
();
requestSpanDisruptor
.
shutdown
();
}
...
...
skywalking-storage-center/skywalking-routing/src/main/java/com/a/eye/skywalking/routing/router/Router.java
浏览文件 @
fdaa3186
...
...
@@ -19,7 +19,6 @@ public class Router implements NodeChangesListener {
private
static
ILog
logger
=
LogManager
.
getLogger
(
Router
.
class
);
private
SpanDisruptor
[]
disruptors
=
new
SpanDisruptor
[
0
];
private
NoopSpanDisruptor
noopSpanPool
=
new
NoopSpanDisruptor
();
public
SpanDisruptor
lookup
(
RequestSpan
requestSpan
)
{
return
getSpanDisruptor
(
requestSpan
.
getRouteKey
());
...
...
@@ -31,7 +30,7 @@ public class Router implements NodeChangesListener {
private
SpanDisruptor
getSpanDisruptor
(
long
routKey
)
{
if
(
disruptors
.
length
==
0
)
{
return
noopSpanPool
;
return
NoopSpanDisruptor
.
INSTANCE
;
}
while
(
true
)
{
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录