Skip to content
Projects
Groups
Snippets
Help
Loading...
Help
Support
Keyboard shortcuts
?
Submit feedback
Contribute to GitLab
Sign in
Toggle navigation
A
app-collection
Project overview
Project overview
Details
Activity
Releases
Repository
Repository
Files
Commits
Branches
Tags
Contributors
Graph
Compare
Issues
0
Issues
0
List
Boards
Labels
Milestones
Merge Requests
0
Merge Requests
0
CI / CD
CI / CD
Pipelines
Jobs
Schedules
Analytics
Analytics
CI / CD
Repository
Value Stream
Wiki
Wiki
Snippets
Snippets
Members
Members
Collapse sidebar
Close sidebar
Activity
Graph
Create a new issue
Jobs
Commits
Issue Boards
Open sidebar
徐顺
app-collection
Commits
37848653
Commit
37848653
authored
Jan 26, 2021
by
RingEric
Browse files
Options
Browse Files
Download
Email Patches
Plain Diff
kafka消费者Api
parent
a49daf90
Changes
1
Hide whitespace changes
Inline
Side-by-side
Showing
1 changed file
with
29 additions
and
0 deletions
+29
-0
src/main/scala/qm/kafka/KafkaConsumer.scala
src/main/scala/qm/kafka/KafkaConsumer.scala
+29
-0
No files found.
src/main/scala/qm/kafka/KafkaConsumer.scala
0 → 100644
View file @
37848653
package
qm.kafka
import
java.util.Properties
import
org.apache.flink.api.common.serialization.SimpleStringSchema
import
org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumer
import
org.apache.kafka.clients.consumer.ConsumerConfig
/**
* @ClassName: KafkaConsumer
* @Description: TODO
* @Create by: LinYoung
* @Date: 2021/1/13 14:54
*/
object
KafkaConsumer
{
def
getConsumer
:
FlinkKafkaConsumer
[
String
]
=
{
val
properties
=
new
Properties
()
// 对接kafka
properties
.
put
(
ConsumerConfig
.
BOOTSTRAP_SERVERS_CONFIG
,
"8.135.22.177:9092"
)
properties
.
put
(
ConsumerConfig
.
GROUP_ID_CONFIG
,
"goods-group"
)
properties
.
put
(
ConsumerConfig
.
KEY_DESERIALIZER_CLASS_CONFIG
,
"org.apache.kafka.common.serialization.StringDeserializer"
)
properties
.
put
(
ConsumerConfig
.
VALUE_DESERIALIZER_CLASS_CONFIG
,
"org.apache.kafka.common.serialization.StringDeserializer"
)
properties
.
put
(
ConsumerConfig
.
AUTO_COMMIT_INTERVAL_MS_CONFIG
,
"1000"
)
properties
.
put
(
ConsumerConfig
.
ENABLE_AUTO_COMMIT_CONFIG
,
"true"
)
val
topic
=
"user_actions"
new
FlinkKafkaConsumer
[
String
](
topic
,
new
SimpleStringSchema
(),
properties
)
}
}
Write
Preview
Markdown
is supported
0%
Try again
or
attach a new file
Attach a file
Cancel
You are about to add
0
people
to the discussion. Proceed with caution.
Finish editing this message first!
Cancel
Please
register
or
sign in
to comment