- ubuntu12.04环境下使用kvm ioctl接口实现最简单的虚拟机
- Ubuntu 通过无线网络安装Ubuntu Server启动系统后连接无线网络的方法
- 在Ubuntu上搭建网桥的方法
- ubuntu 虚拟机上网方式及相关配置详解
CFSDN坚持开源创造价值,我们致力于搭建一个资源共享平台,让每一个IT人在这里找到属于你的精彩世界.
这篇CFSDN的博客文章Java实现Kafka生产者和消费者的示例由作者收集整理,如果你对这篇文章有兴趣,记得点赞哟.
Kafka简介 。
Kafka是由Apache软件基金会开发的一个开源流处理平台,由Scala和Java编写。Kafka的目标是为处理实时数据提供一个统1、高吞吐、低延迟的平台.
方式一:kafka-clients 。
引入依赖 。
在pom.xml文件中,引入kafka-clients依赖:
1
2
3
4
5
|
<
dependency
>
<
groupId
>org.apache.kafka</
groupId
>
<
artifactId
>kafka-clients</
artifactId
>
<
version
>2.3.1</
version
>
</
dependency
>
|
生产者 。
创建一个KafkaProducer的生产者实例:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
|
@Configuration
public
class
Config {
public
final
static
String bootstrapServers =
"127.0.0.1:9092"
;
@Bean
(destroyMethod =
"close"
)
public
KafkaProducer<String, String> kafkaProducer() {
Properties props =
new
Properties();
//设置Kafka服务器地址
props.put(
"bootstrap.servers"
, bootstrapServers);
//设置数据key的序列化处理类
props.put(
"key.serializer"
, StringSerializer.
class
.getName());
//设置数据value的序列化处理类
props.put(
"value.serializer"
, StringSerializer.
class
.getName());
KafkaProducer<String, String> producer =
new
KafkaProducer<>(props);
return
producer;
}
}
|
在Controller中进行使用:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
|
@RestController
@Slf4j
public
class
Controller {
@Autowired
private
KafkaProducer<String, String> kafkaProducer;
@RequestMapping
(
"/kafkaClientsSend"
)
public
String send() {
String uuid = UUID.randomUUID().toString();
RecordMetadata recordMetadata =
null
;
try
{
//将消息发送到Kafka服务器的名称为“one-more-topic”的Topic中
recordMetadata = kafkaProducer.send(
new
ProducerRecord<>(
"one-more-topic"
, uuid)).get();
log.info(
"recordMetadata: {}"
, recordMetadata);
log.info(
"uuid: {}"
, uuid);
}
catch
(Exception e) {
log.error(
"send fail, uuid: {}"
, uuid, e);
}
return
uuid;
}
}
|
消费者 。
创建一个KafkaConsumer的消费者实例:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
|
@Configuration
public
class
Config {
public
final
static
String groupId =
"kafka-clients-group"
;
public
final
static
String bootstrapServers =
"127.0.0.1:9092"
;
@Bean
(destroyMethod =
"close"
)
public
KafkaConsumer<String, String> kafkaConsumer() {
Properties props =
new
Properties();
//设置Kafka服务器地址
props.put(
"bootstrap.servers"
, bootstrapServers);
//设置消费组
props.put(
"group.id"
, groupId);
//设置数据key的反序列化处理类
props.put(
"key.deserializer"
, StringDeserializer.
class
.getName());
//设置数据value的反序列化处理类
props.put(
"value.deserializer"
, StringDeserializer.
class
.getName());
props.put(
"enable.auto.commit"
,
"true"
);
props.put(
"auto.commit.interval.ms"
,
"1000"
);
props.put(
"session.timeout.ms"
,
"30000"
);
KafkaConsumer<String, String> kafkaConsumer =
new
KafkaConsumer<>(props);
//订阅名称为“one-more-topic”的Topic的消息
kafkaConsumer.subscribe(Arrays.asList(
"one-more-topic"
));
return
kafkaConsumer;
}
}
|
在Controller中进行使用:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
|
@RestController
@Slf4j
public
class
Controller {
@Autowired
private
KafkaConsumer<String, String> kafkaConsumer;
@RequestMapping
(
"/receive"
)
public
List<String> receive() {
从Kafka服务器中的名称为“one-more-topic”的Topic中消费消息
ConsumerRecords<String, String> records = kafkaConsumer.poll(Duration.ofSeconds(
1
));
List<String> messages =
new
ArrayList<>(records.count());
for
(ConsumerRecord<String, String> record : records.records(
"one-more-topic"
)) {
String message = record.value();
log.info(
"message: {}"
, message);
messages.add(message);
}
return
messages;
}
}
|
。
方式二:spring-kafka 。
使用kafka-clients需要我们自己创建生产者或者消费者的bean,如果我们的项目基于SpringBoot构建,那么使用spring-kafka就方便多了.
引入依赖 。
在pom.xml文件中,引入spring-kafka依赖:
1
2
3
4
5
|
<
dependency
>
<
groupId
>org.springframework.kafka</
groupId
>
<
artifactId
>spring-kafka</
artifactId
>
<
version
>2.3.12.RELEASE</
version
>
</
dependency
>
|
生产者 。
在application.yml文件中增加配置:
1
2
3
4
5
6
7
|
spring:
kafka:
#Kafka服务器地址
bootstrap-servers: 127.0.0.1:9092
producer:
#设置数据value的序列化处理类
value-serializer: org.apache.kafka.common.serialization.StringSerializer
|
在Controller中注入KafkaTemplate就可以直接使用了,代码如下:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
|
@RestController
@Slf4j
public
class
Controller {
@Autowired
private
KafkaTemplate<String, String> template;
@RequestMapping
(
"/springKafkaSend"
)
public
String send() {
String uuid = UUID.randomUUID().toString();
//将消息发送到Kafka服务器的名称为“one-more-topic”的Topic中
this
.template.send(
"one-more-topic"
, uuid);
log.info(
"uuid: {}"
, uuid);
return
uuid;
}
}
|
消费者 。
在application.yml文件中增加配置:
1
2
3
4
5
6
7
|
spring:
kafka:
#Kafka服务器地址
bootstrap-servers: 127.0.0.1:9092
consumer:
#设置数据value的反序列化处理类
value-deserializer: org.apache.kafka.common.serialization.StringDeserializer
|
创建一个可以被Spring框架扫描到的类,并且在方法上加上@KafkaListener注解,就可以消费消息了,代码如下:
1
2
3
4
5
6
7
8
9
10
11
12
13
|
@Component
@Slf4j
public
class
Receiver {
@KafkaListener
(topics =
"one-more-topic"
, groupId =
"spring-kafka-group"
)
public
void
listen(ConsumerRecord<?, ?> record) {
Optional<?> kafkaMessage = Optional.ofNullable(record.value());
if
(kafkaMessage.isPresent()) {
String message = (String) kafkaMessage.get();
log.info(
"message: {}"
, message);
}
}
}
|
到此这篇关于Java实现Kafka生产者和消费者的示例的文章就介绍到这了,更多相关Java Kafka生产者和消费者 内容请搜索我以前的文章或继续浏览下面的相关文章希望大家以后多多支持我! 。
原文链接:https://blog.csdn.net/heihaozi/article/details/111042472 。
最后此篇关于Java实现Kafka生产者和消费者的示例的文章就讲到这里了,如果你想了解更多关于Java实现Kafka生产者和消费者的示例的内容请搜索CFSDN的文章或继续浏览相关文章,希望大家以后支持我的博客! 。
我在 Windows 机器上启动 Kafka-Server 时出现以下错误。我已经从以下链接下载了 Scala 2.11 - kafka_2.11-2.1.0.tgz:https://kafka.ap
关于Apache-Kafka messaging queue . 我已经从 Kafka 下载页面下载了 Apache Kafka。我已将其提取到 /opt/apache/installed/kafka
假设我有 Kafka 主题 cars。 我还有一个消费者组 cars-consumers 订阅了 cars 主题。 cars-consumers 消费者组当前位于偏移量 89。 当我现在删除 cars
我想知道什么最适合我:Kafka 流或 Kafka 消费者 api 或 Kafka 连接? 我想从主题中读取数据,然后进行一些处理并写入数据库。所以我编写了消费者,但我觉得我可以编写 Kafka 流应
我曾研究过一些 Kafka 流应用程序和 Kafka 消费者应用程序。最后,Kafka流不过是消费来自Kafka的实时事件的消费者。因此,我无法弄清楚何时使用 Kafka 流或为什么我们应该使用
Kafka Acknowledgement 和 Kafka 消费者 commitSync() 有什么区别 两者都用于手动偏移管理,并希望两者同步工作。 请协助 最佳答案 使用 spring-kafka
如何在 Kafka 代理上代理 Apache Kafka 生产者请求,并重定向到单独的 Kafka 集群? 在我的特定情况下,无法更新写入此集群的客户端。这意味着,执行以下操作是不可行的: 更新客户端
我需要在 Kafka 10 中命名我的消费者,就像我在 Kafka 8 中所做的一样,因为我有脚本可以嗅出并进一步使用这些信息。 显然,consumer.id 的默认命名已更改(并且现在还单独显示了
1.概述 我们会看到zk的数据中有一个节点/log_dir_event_notification/,这是一个序列号持久节点 这个节点在kafka中承担的作用是: 当某个Broker上的LogDir出现
我正在使用以下命令: bin/kafka-console-producer.sh --broker-list localhost:9092 --topic test.topic --property
我很难理解 Java Spring Boot 中的一些 Kafka 概念。我想针对在服务器上运行的真实 Kafka 代理测试消费者,该服务器有一些生产者已将数据写入/已经将数据写入各种主题。我想与服务
我的场景是我使用了很多共享前缀的 Kafka 主题(例如 house.door, house.room ) 并使用 Kafka 流正则表达式主题模式 API 使用所有主题。 一切看起来都不错,我得到了
有没有办法以编程方式获取kafka集群的版本?例如,使用AdminClient应用程序接口(interface)。 我想在消费者/生产者应用程序中识别 kafka 集群的版本。 最佳答案 目前无法检索
每当我尝试重新启动 kafka 时,它都会出现以下错误。一旦我删除/tmp/kafka-logs 它就会得到解决,但它也会删除我的主题。 有办法解决吗? ERROR Error while
我是 Apache Kafka 的新用户,我仍在了解内部结构。 在我的用例中,我需要从 Kafka Producer 客户端动态增加主题的分区数。 我发现了其他类似的 questions关于增加分区大
正如 Kafka 文档所示,一种方法是通过 kafka.tools.MirrorMaker 来实现这一点。但是,我需要将一个主题(比如 测试 带 1 个分区)(其内容和元数据)从生产环境复制到没有连接
我已经在集群中配置了 3 个 kafka,我正在尝试与 spring-kafka 一起使用。 但是在我杀死 kafka 领导者之后,我无法将其他消息发送到队列中。 我将 spring.kafka.bo
我的 kafka sink 连接器从多个主题(配置了 10 个任务)读取,并处理来自所有主题的 300 条记录。根据每个记录中保存的信息,连接器可以执行某些操作。 以下是触发器记录中键值对的示例: "
我有以下 kafka 流代码 public class KafkaStreamHandler implements Processor{ private ProcessorConte
当 kafka-streams 应用程序正在运行并且 Kafka 突然关闭时,应用程序进入“等待”模式,发送警告日志的消费者和生产者线程无法连接,当 Kafka 回来时,一切都应该(理论上)去恢复正常
我是一名优秀的程序员,十分优秀!