【发布时间】:2018-04-30 01:57:47
【问题描述】:
我一直在尝试让入站 SubscribableChannel 和出站 MessageChannel 在我的 Spring Boot 应用程序中工作。
我已经成功设置了kafka通道并测试成功。
此外,我创建了一个基本的 Spring Boot 应用程序,用于测试从通道添加和接收内容。
我遇到的问题是,当我将等效代码放入它所属的应用程序时,消息似乎永远不会被发送或接收。通过调试很难确定发生了什么,但对我来说唯一看起来不同的是频道名称。在工作的 impl 中,通道名称就像在非工作应用程序中的 application.channel 其 localhost:8080/channel。
我想知道是否有一些 Spring Boot 配置阻止或将通道的创建更改为不同的通道源?
有人遇到过类似的问题吗?
application.yml
spring:
datasource:
url: jdbc:h2:mem:dpemail;DB_CLOSE_DELAY=-1;DB_CLOSE_ON_EXIT=FALSE
platform: h2
username: hello
password:
driverClassName: org.h2.Driver
jpa:
properties:
hibernate:
show_sql: true
use_sql_comments: true
format_sql: true
cloud:
stream:
kafka:
binder:
brokers: localhost:9092
bindings:
email-in:
destination: email
contentType: application/json
email-out:
destination: email
contentType: application/json
电子邮件
public class Email {
private long timestamp;
private String message;
public long getTimestamp() {
return timestamp;
}
public void setTimestamp(long timestamp) {
this.timestamp = timestamp;
}
public String getMessage() {
return message;
}
public void setMessage(String message) {
this.message = message;
}
}
绑定配置
@EnableBinding(EmailQueues.class)
public class EmailQueueConfiguration {
}
界面
public interface EmailQueues {
String INPUT = "email-in";
String OUTPUT = "email-out";
@Input(INPUT)
SubscribableChannel inboundEmails();
@Output(OUTPUT)
MessageChannel outboundEmails();
}
控制器
@RestController
@RequestMapping("/queue")
public class EmailQueueController {
private EmailQueues emailQueues;
@Autowired
public EmailQueueController(EmailQueues emailQueues) {
this.emailQueues = emailQueues;
}
@RequestMapping(value = "sendEmail", method = POST)
@ResponseStatus(ACCEPTED)
public void sendToQueue() {
MessageChannel messageChannel = emailQueues.outboundEmails();
Email email = new Email();
email.setMessage("hello world: " + System.currentTimeMillis());
email.setTimestamp(System.currentTimeMillis());
messageChannel.send(MessageBuilder.withPayload(email).setHeader(MessageHeaders.CONTENT_TYPE, MimeTypeUtils.APPLICATION_JSON).build());
}
@StreamListener(EmailQueues.INPUT)
public void handleEmail(@Payload Email email) {
System.out.println("received: " + email.getMessage());
}
}
我不确定使用 Spring-Cloud、Spring-Cloud-Sleuth 的继承配置项目之一是否会阻止它工作,但即使我删除它仍然没有。但与我的应用程序可以使用上述代码不同,我从未看到正在配置 ConsumeConfig,例如:
o.a.k.clients.consumer.ConsumerConfig : ConsumerConfig values:
auto.commit.interval.ms = 100
auto.offset.reset = latest
bootstrap.servers = [localhost:9092]
check.crcs = true
client.id = consumer-2
connections.max.idle.ms = 540000
enable.auto.commit = false
exclude.internal.topics = true
(此配置是我在运行上述代码时在我的基本 Spring Boot 应用程序中看到的,并且代码可以从 kafka 通道写入和读取)......
我假设我正在使用的一个库中存在一些 over spring boot 配置,创建了一种不同类型的通道,我只是找不到该配置是什么。
【问题讨论】:
-
考虑粘贴您的代码
-
任何人都可以帮助您的唯一方法是发布您的代码和配置 - 编辑问题,不要将代码等放入 cmets。
-
代码添加到原始帖子。
标签: spring-boot apache-kafka spring-cloud-stream