Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
doujutun3207
flink
提交
123be227
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,发现更多精彩内容 >>
提交
123be227
编写于
6月 28, 2016
作者:
R
Robert Metzger
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
[FLINK-4085][Kinesis] Set Flink-specific user agent
This closes #2175
上级
256c9c4d
变更
1
隐藏空白更改
内联
并排
Showing
1 changed file
with
11 addition
and
6 deletion
+11
-6
flink-streaming-connectors/flink-connector-kinesis/src/main/java/org/apache/flink/streaming/connectors/kinesis/proxy/KinesisProxy.java
...link/streaming/connectors/kinesis/proxy/KinesisProxy.java
+11
-6
未找到文件。
flink-streaming-connectors/flink-connector-kinesis/src/main/java/org/apache/flink/streaming/connectors/kinesis/proxy/KinesisProxy.java
浏览文件 @
123be227
...
...
@@ -17,6 +17,8 @@
package
org.apache.flink.streaming.connectors.kinesis.proxy
;
import
com.amazonaws.ClientConfiguration
;
import
com.amazonaws.ClientConfigurationFactory
;
import
com.amazonaws.regions.Region
;
import
com.amazonaws.regions.Regions
;
import
com.amazonaws.services.kinesis.AmazonKinesisClient
;
...
...
@@ -29,6 +31,7 @@ import com.amazonaws.services.kinesis.model.LimitExceededException;
import
com.amazonaws.services.kinesis.model.ResourceNotFoundException
;
import
com.amazonaws.services.kinesis.model.StreamStatus
;
import
com.amazonaws.services.kinesis.model.Shard
;
import
org.apache.flink.runtime.util.EnvironmentInformation
;
import
org.apache.flink.streaming.connectors.kinesis.config.KinesisConfigConstants
;
import
org.apache.flink.streaming.connectors.kinesis.model.KinesisStreamShard
;
import
org.apache.flink.streaming.connectors.kinesis.util.AWSUtil
;
...
...
@@ -58,9 +61,6 @@ public class KinesisProxy {
/** The actual Kinesis client from the AWS SDK that we will be using to make calls */
private
final
AmazonKinesisClient
kinesisClient
;
/** The AWS region that this proxy will be making calls to */
private
final
String
regionId
;
/** Configuration properties of this Flink Kinesis Connector */
private
final
Properties
configProps
;
...
...
@@ -72,9 +72,14 @@ public class KinesisProxy {
public
KinesisProxy
(
Properties
configProps
)
{
this
.
configProps
=
checkNotNull
(
configProps
);
this
.
regionId
=
configProps
.
getProperty
(
KinesisConfigConstants
.
CONFIG_AWS_REGION
);
AmazonKinesisClient
client
=
new
AmazonKinesisClient
(
AWSUtil
.
getCredentialsProvider
(
configProps
).
getCredentials
());
client
.
setRegion
(
Region
.
getRegion
(
Regions
.
fromName
(
this
.
regionId
)));
/* The AWS region that this proxy will be making calls to */
String
regionId
=
configProps
.
getProperty
(
KinesisConfigConstants
.
CONFIG_AWS_REGION
);
// set Flink as a user agent
ClientConfiguration
config
=
new
ClientConfigurationFactory
().
getConfig
();
config
.
setUserAgent
(
"Apache Flink "
+
EnvironmentInformation
.
getVersion
()
+
" ("
+
EnvironmentInformation
.
getRevisionInformation
().
commitId
+
") Kinesis Connector"
);
AmazonKinesisClient
client
=
new
AmazonKinesisClient
(
AWSUtil
.
getCredentialsProvider
(
configProps
).
getCredentials
(),
config
);
client
.
setRegion
(
Region
.
getRegion
(
Regions
.
fromName
(
regionId
)));
this
.
kinesisClient
=
client
;
}
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录