【发布时间】:2018-09-18 04:16:56
【问题描述】:
我正在尝试创建 kstreams bean 并在我的服务中自动装配它。但是,即使我得到相同的对象 stream.print() 也没有给出任何价值,但在同一个 bean 中的 print 正在工作。我想我没有得到 Same StreamBuilder 的配置。
配置文件
@Configuration
@EnableKafka
@EnableKafkaStreams
public class KafkaStreamsConfiguration {
@Autowired private KafkaProperties kafkaProperties;
@Bean(name = KafkaStreamsDefaultConfiguration.DEFAULT_STREAMS_CONFIG_BEAN_NAME)
public StreamsConfig kStreamsConfigs() {
Map<String, Object> props = new HashMap<>();
props.put(StreamsConfig.APPLICATION_ID_CONFIG, "test-streams2");
props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaProperties.getBootstrapServers());
props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass());
props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass());
props.put(JsonDeserializer.DEFAULT_KEY_TYPE, String.class);
props.put(JsonDeserializer.DEFAULT_VALUE_TYPE, String.class);
return new StreamsConfig(props);
}
@Bean
public KStream<String, String> kStreamJson(StreamsBuilder builder) {
KStream<String, String> stream = builder.stream("topictest", Consumed.with(Serdes.String(), Serdes.String()));
//stream.print();
return stream;
}
}
服务
这里的打印函数没有抛出任何错误,也没有打印任何值
@Service
public class KStreamsService {
@Autowired
KStream<String, String> kStream;
void process() {
System.out.println("Hai");
kStream.print();
}
}
主要
@SpringBootApplication
public class KStreamsApplication {
@Autowired
KStreamsService kStreamsService;
public static void main(String[] args) {
SpringApplication.run(KStreamsApplication.class, args);
}
private void run() {
kStreamsService.process();
}
}
我在这里做错了吗?
【问题讨论】:
-
请注意,
print()可能会缓冲数据(参见issues.apache.org/jira/browse/KAFKA-7326)。解决方法是foreach()并在您的用户代码中调用System.out.println()。
标签: apache-kafka apache-kafka-streams spring-kafka