Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
jobily
Questdb
提交
888cef74
Q
Questdb
项目概览
jobily
/
Questdb
11 个月 前同步成功
通知
1
Star
0
Fork
0
代码
文件
提交
分支
Tags
贡献者
分支图
Diff
Issue
0
列表
看板
标记
里程碑
合并请求
0
DevOps
流水线
流水线任务
计划
Wiki
0
Wiki
分析
仓库
DevOps
项目成员
Pages
Q
Questdb
项目概览
项目概览
详情
发布
仓库
仓库
文件
提交
分支
标签
贡献者
分支图
比较
Issue
0
Issue
0
列表
看板
标记
里程碑
合并请求
0
合并请求
0
Pages
DevOps
DevOps
流水线
流水线任务
计划
分析
分析
仓库分析
DevOps
Wiki
0
Wiki
成员
成员
收起侧边栏
关闭侧边栏
动态
分支图
创建新Issue
流水线任务
提交
Issue看板
体验新版 GitCode,发现更多精彩内容 >>
未验证
提交
888cef74
编写于
10月 21, 2021
作者:
A
Alex Pelagenko
提交者:
GitHub
10月 21, 2021
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
fix(core): fix SO in concurrent FanOut.and (#1460)
上级
826f3e67
变更
2
隐藏空白更改
内联
并排
Showing
2 changed file
with
40 addition
and
1 deletion
+40
-1
core/src/main/java/io/questdb/mp/FanOut.java
core/src/main/java/io/questdb/mp/FanOut.java
+4
-1
core/src/test/java/io/questdb/mp/ConcurrentTest.java
core/src/test/java/io/questdb/mp/ConcurrentTest.java
+36
-0
未找到文件。
core/src/main/java/io/questdb/mp/FanOut.java
浏览文件 @
888cef74
...
@@ -61,6 +61,8 @@ public class FanOut implements Barrier {
...
@@ -61,6 +61,8 @@ public class FanOut implements Barrier {
public
FanOut
and
(
Barrier
barrier
)
{
public
FanOut
and
(
Barrier
barrier
)
{
Holder
_new
;
Holder
_new
;
boolean
barrierNotSetUp
=
true
;
do
{
do
{
Holder
h
=
this
.
holder
;
Holder
h
=
this
.
holder
;
// read barrier to make sure "holder" read doesn't fall below this
// read barrier to make sure "holder" read doesn't fall below this
...
@@ -68,8 +70,9 @@ public class FanOut implements Barrier {
...
@@ -68,8 +70,9 @@ public class FanOut implements Barrier {
if
(
h
.
barriers
.
indexOf
(
barrier
)
>
-
1
)
{
if
(
h
.
barriers
.
indexOf
(
barrier
)
>
-
1
)
{
return
this
;
return
this
;
}
}
if
(
this
.
barrier
!=
null
)
{
if
(
barrierNotSetUp
&&
this
.
barrier
!=
null
)
{
barrier
.
root
().
setBarrier
(
this
.
barrier
);
barrier
.
root
().
setBarrier
(
this
.
barrier
);
barrierNotSetUp
=
false
;
}
}
_new
=
new
Holder
();
_new
=
new
Holder
();
_new
.
barriers
.
addAll
(
h
.
barriers
);
_new
.
barriers
.
addAll
(
h
.
barriers
);
...
...
core/src/test/java/io/questdb/mp/ConcurrentTest.java
浏览文件 @
888cef74
...
@@ -35,6 +35,7 @@ import java.util.Arrays;
...
@@ -35,6 +35,7 @@ import java.util.Arrays;
import
java.util.concurrent.BrokenBarrierException
;
import
java.util.concurrent.BrokenBarrierException
;
import
java.util.concurrent.CountDownLatch
;
import
java.util.concurrent.CountDownLatch
;
import
java.util.concurrent.CyclicBarrier
;
import
java.util.concurrent.CyclicBarrier
;
import
java.util.concurrent.atomic.AtomicInteger
;
import
java.util.concurrent.locks.LockSupport
;
import
java.util.concurrent.locks.LockSupport
;
public
class
ConcurrentTest
{
public
class
ConcurrentTest
{
...
@@ -389,6 +390,41 @@ public class ConcurrentTest {
...
@@ -389,6 +390,41 @@ public class ConcurrentTest {
}
}
}
}
@Test
public
void
testConcurrentFanOutAnd
()
{
int
cycle
=
1024
;
SPSequence
pubSeq
=
new
SPSequence
(
cycle
);
FanOut
fout
=
new
FanOut
();
pubSeq
.
then
(
fout
).
then
(
pubSeq
);
int
threads
=
2
;
CyclicBarrier
start
=
new
CyclicBarrier
(
threads
);
SOCountDownLatch
latch
=
new
SOCountDownLatch
(
threads
);
int
iterations
=
30
;
AtomicInteger
doneCount
=
new
AtomicInteger
();
for
(
int
i
=
0
;
i
<
threads
;
i
++)
{
new
Thread
(()
->
{
try
{
start
.
await
();
for
(
int
j
=
0
;
j
<
iterations
;
j
++)
{
SCSequence
consumer
=
new
SCSequence
();
FanOut
fout2
=
fout
.
and
(
consumer
);
fout2
.
remove
(
consumer
);
}
doneCount
.
addAndGet
(
iterations
);
}
catch
(
InterruptedException
|
BrokenBarrierException
e
)
{
e
.
printStackTrace
();
}
finally
{
latch
.
countDown
();
}
}).
start
();
}
latch
.
await
();
Assert
.
assertEquals
(
threads
*
iterations
,
doneCount
.
get
());
}
static
void
publishEOE
(
RingQueue
<
Event
>
queue
,
Sequence
sequence
)
{
static
void
publishEOE
(
RingQueue
<
Event
>
queue
,
Sequence
sequence
)
{
long
cursor
=
sequence
.
nextBully
();
long
cursor
=
sequence
.
nextBully
();
queue
.
get
(
cursor
).
value
=
Integer
.
MIN_VALUE
;
queue
.
get
(
cursor
).
value
=
Integer
.
MIN_VALUE
;
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录