Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
doujutun3207
flink
提交
12197b3b
F
flink
项目概览
doujutun3207
/
flink
与 Fork 源项目一致
从无法访问的项目Fork
通知
24
Star
0
Fork
0
代码
文件
提交
分支
Tags
贡献者
分支图
Diff
Issue
0
列表
看板
标记
里程碑
合并请求
0
Wiki
0
Wiki
分析
仓库
DevOps
项目成员
Pages
F
flink
项目概览
项目概览
详情
发布
仓库
仓库
文件
提交
分支
标签
贡献者
分支图
比较
Issue
0
Issue
0
列表
看板
标记
里程碑
合并请求
0
合并请求
0
Pages
分析
分析
仓库分析
DevOps
Wiki
0
Wiki
成员
成员
收起侧边栏
关闭侧边栏
动态
分支图
创建新Issue
提交
Issue看板
体验新版 GitCode,发现更多精彩内容 >>
未验证
提交
12197b3b
编写于
4月 20, 2020
作者:
A
Aljoscha Krettek
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
[minor] Fix warnings in AbstractFetcher(Test)
上级
8e5685d3
变更
2
隐藏空白更改
内联
并排
Showing
2 changed file
with
14 addition
and
10 deletion
+14
-10
flink-connectors/flink-connector-kafka-base/src/main/java/org/apache/flink/streaming/connectors/kafka/internals/AbstractFetcher.java
...streaming/connectors/kafka/internals/AbstractFetcher.java
+5
-3
flink-connectors/flink-connector-kafka-base/src/test/java/org/apache/flink/streaming/connectors/kafka/internals/AbstractFetcherTest.java
...aming/connectors/kafka/internals/AbstractFetcherTest.java
+9
-7
未找到文件。
flink-connectors/flink-connector-kafka-base/src/main/java/org/apache/flink/streaming/connectors/kafka/internals/AbstractFetcher.java
浏览文件 @
12197b3b
...
...
@@ -130,9 +130,11 @@ public abstract class AbstractFetcher<T, KPH> {
*/
private
final
MetricGroup
consumerMetricGroup
;
@SuppressWarnings
(
"DeprecatedIsStillUsed"
)
@Deprecated
private
final
MetricGroup
legacyCurrentOffsetsMetricGroup
;
@SuppressWarnings
(
"DeprecatedIsStillUsed"
)
@Deprecated
private
final
MetricGroup
legacyCommittedOffsetsMetricGroup
;
...
...
@@ -185,7 +187,7 @@ public abstract class AbstractFetcher<T, KPH> {
userCodeClassLoader
);
// check that all seed partition states have a defined offset
for
(
KafkaTopicPartitionState
partitionState
:
subscribedPartitionStates
)
{
for
(
KafkaTopicPartitionState
<?>
partitionState
:
subscribedPartitionStates
)
{
if
(!
partitionState
.
isOffsetDefined
())
{
throw
new
IllegalArgumentException
(
"The fetcher was assigned seed partitions with undefined initial offsets."
);
}
...
...
@@ -204,7 +206,7 @@ public abstract class AbstractFetcher<T, KPH> {
// if we have periodic watermarks, kick off the interval scheduler
if
(
timestampWatermarkMode
==
PERIODIC_WATERMARKS
)
{
@SuppressWarnings
(
"unchecked"
)
PeriodicWatermarkEmitter
periodicEmitter
=
new
PeriodicWatermarkEmitter
(
PeriodicWatermarkEmitter
<
KPH
>
periodicEmitter
=
new
PeriodicWatermarkEmitter
(
subscribedPartitionStates
,
sourceContext
,
processingTimeProvider
,
...
...
@@ -303,7 +305,7 @@ public abstract class AbstractFetcher<T, KPH> {
return
offsets
.
entrySet
()
.
stream
()
.
filter
(
entry
->
!
KafkaTopicPartitionStateSentinel
.
isSentinel
(
entry
.
getValue
()))
.
collect
(
Collectors
.
toMap
(
entry
->
entry
.
getKey
(),
entry
->
entry
.
getValue
()
));
.
collect
(
Collectors
.
toMap
(
Map
.
Entry
::
getKey
,
Map
.
Entry
::
getValue
));
}
/**
...
...
flink-connectors/flink-connector-kafka-base/src/test/java/org/apache/flink/streaming/connectors/kafka/internals/AbstractFetcherTest.java
浏览文件 @
12197b3b
...
...
@@ -443,7 +443,7 @@ public class AbstractFetcherTest {
// ------------------------------------------------------------------------
private
static
final
class
TestFetcher
<
T
>
extends
AbstractFetcher
<
T
,
Object
>
{
Optional
<
Map
<
KafkaTopicPartition
,
Long
>>
lastCommittedOffsets
=
Optional
.
empty
()
;
Map
<
KafkaTopicPartition
,
Long
>
lastCommittedOffsets
=
null
;
private
final
OneShotLatch
fetchLoopWaitLatch
;
private
final
OneShotLatch
stateIterationBlockLatch
;
...
...
@@ -499,7 +499,7 @@ public class AbstractFetcherTest {
@Override
public
void
runFetchLoop
()
throws
Exception
{
if
(
fetchLoopWaitLatch
!=
null
)
{
for
(
KafkaTopicPartitionState
ignored
:
subscribedPartitionStates
())
{
for
(
KafkaTopicPartitionState
<?>
ignored
:
subscribedPartitionStates
())
{
fetchLoopWaitLatch
.
trigger
();
stateIterationBlockLatch
.
await
();
}
...
...
@@ -521,13 +521,13 @@ public class AbstractFetcherTest {
@Override
protected
void
doCommitInternalOffsetsToKafka
(
Map
<
KafkaTopicPartition
,
Long
>
offsets
,
@Nonnull
KafkaCommitCallback
callback
)
throws
Exception
{
lastCommittedOffsets
=
Optional
.
of
(
offsets
)
;
@Nonnull
KafkaCommitCallback
callback
)
{
lastCommittedOffsets
=
offsets
;
callback
.
onSuccess
();
}
public
Optional
<
Map
<
KafkaTopicPartition
,
Long
>>
getLastCommittedOffsets
()
{
return
lastCommittedOffsets
;
return
Optional
.
ofNullable
(
lastCommittedOffsets
)
;
}
}
...
...
@@ -537,7 +537,7 @@ public class AbstractFetcherTest {
AbstractFetcher
<
T
,
KPH
>
fetcher
,
T
record
,
KafkaTopicPartitionState
<
KPH
>
partitionState
,
long
offset
)
throws
Exception
{
long
offset
)
{
ArrayDeque
<
T
>
recordQueue
=
new
ArrayDeque
<>();
recordQueue
.
add
(
record
);
...
...
@@ -552,7 +552,7 @@ public class AbstractFetcherTest {
AbstractFetcher
<
T
,
KPH
>
fetcher
,
List
<
T
>
records
,
KafkaTopicPartitionState
<
KPH
>
partitionState
,
long
offset
)
throws
Exception
{
long
offset
)
{
ArrayDeque
<
T
>
recordQueue
=
new
ArrayDeque
<>(
records
);
fetcher
.
emitRecordsWithTimestamps
(
...
...
@@ -566,6 +566,7 @@ public class AbstractFetcherTest {
return
new
ArrayDeque
<>();
}
@SuppressWarnings
(
"deprecation"
)
private
static
class
PeriodicTestExtractor
implements
AssignerWithPeriodicWatermarks
<
Long
>
{
private
volatile
long
maxTimestamp
=
Long
.
MIN_VALUE
;
...
...
@@ -583,6 +584,7 @@ public class AbstractFetcherTest {
}
}
@SuppressWarnings
(
"deprecation"
)
private
static
class
PunctuatedTestExtractor
implements
AssignerWithPunctuatedWatermarks
<
Long
>
{
@Override
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录