我正在尝试在我的 Windows 中安装 zookeeper。无论我在zookeeper + Kafka - Unable to create data directory 中遵循哪个建议,我都会收到以下错误。
我以管理员身份运行它,并尝试了所有这些选项:
#dataDir=/tmp/zookeeper
#dataDir=:\zookeeper-3.4.14\
#dataDir=C:\\_d\\WSs\\kafka\\zookeeper-3.4.14\\data
#dataDir=:\\\\zookeeper\\\\data
dataDir=C:\\_d\\WSs\\kafka\\zookeeper-3.4.14
Run Code Online (Sandbox Code Playgroud)
我认为这无关紧要,但让我在这里补充一下:我有 Java 11。
任何想法为什么会发生将不胜感激。
完整日志
C:\Windows\system32>zkserver
C:\Windows\system32>call "C:\Program Files\Java\jdk-11.0.2"\bin\java "-Dzookeeper.log.dir=C:\_d\WSs\kafka\zookeeper-3.4.14\bin\.." "-Dzookeeper.root.logger=INFO,CONSOLE" -cp "C:\_d\WSs\kafka\zookeeper-3.4.14\bin\..\build\classes;C:\_d\WSs\kafka\zookeeper-3.4.14\bin\..\build\lib\*;C:\_d\WSs\kafka\zookeeper-3.4.14\bin\..\*;C:\_d\WSs\kafka\zookeeper-3.4.14\bin\..\lib\*;C:\_d\WSs\kafka\zookeeper-3.4.14\bin\..\conf" org.apache.zookeeper.server.quorum.QuorumPeerMain "C:\_d\WSs\kafka\zookeeper-3.4.14\bin\..\conf\zoo.cfg"
2019-04-18 15:17:42,629 [myid:] - INFO [main:QuorumPeerConfig@136] - Reading configuration from: C:\_d\WSs\kafka\zookeeper-3.4.14\bin\..\conf\zoo.cfg
2019-04-18 15:17:42,644 [myid:] - INFO [main:DatadirCleanupManager@78] - autopurge.snapRetainCount set to 3
2019-04-18 15:17:42,644 [myid:] - INFO [main:DatadirCleanupManager@79] - autopurge.purgeInterval set to 0
2019-04-18 15:17:42,644 [myid:] - INFO [main:DatadirCleanupManager@101] - Purge task is not …Run Code Online (Sandbox Code Playgroud) 我的问题中粘贴的声明是从https://developer.mozilla.org/en-US/docs/Web/Web_Components/Using_custom_elements#Using_the_lifecycle_callbacks复制的。
作为一名没有 WebComponent 经验的开发人员,我试图了解迄今为止推荐的所有经验法则和最佳实践。
继续阅读它说“...使用 Node.isConnected 来确保”。它的含义很明显:检查它是否仍然连接,但至少对我来说,尚不清楚我应该采取什么措施来解决它,或者在某些情况下我应该期待什么。
我的情况是我正在创建一个 Web 组件来侦听 SSE(服务器发送事件)。这对于 alife 仪表板和其他几个场景非常有用。SSE事件在从Kafka Stream消费后基本上将由NodeJs或Spring Webflux响应。
到目前为止,我所做的所有简单示例都没有遇到任何元素在连接回调期间不再连接的问题。
此外,我没有阅读最佳实践中关于“不再连接的元素”的任何建议。
我读了一些精彩的讨论:
从那里我了解到我始终可以信任这个生命周期构造函数 -->connectedCallback -->disconnectedCallback。
和
当所有子自定义元素已连接时如何拥有“connectedCallback”
基本上我了解到没有一个特定的方法“在所有孩子都升级后调用”
这两个问题都接近我的问题,但它没有回答我:应该注意哪些挑战或风险,或者如何解决“一旦元素不再连接,可能会调用connectedCallback”的可能性?在我上述的情况下,我是否缺少任何治疗?我是否应该创建一些观察者,当该元素不再可用时触发,以重新创建事件源对象并再次向此类事件源对象添加侦听器?
我粘贴了下面的代码进行说明,完整的 Web 组件示例可以从https://github.com/jimisdrpc/simplest-webcomponet克隆,其后端可以从https://github.com/jimisdrpc/simplest-kafkaconsumer克隆。
const template = document.createElement('template');
template.innerHTML = `<input id="inputKafka"/> `;
class InputKafka extends HTMLElement {
constructor() {
super();
}
connectedCallback() {
this.attachShadow({mode: 'open'})
this.shadowRoot.appendChild(template.content.cloneNode(true))
const inputKafka = this.shadowRoot.getElementById('inputKafka');
var source = new EventSource('http://localhost:5000/kafka_sse');
source.addEventListener('sendMsgFromKafka', function(e) {
console.log('fromKafka');
inputKafka.value = e.data;
}, false);
}
attributeChangedCallback(name, …Run Code Online (Sandbox Code Playgroud) html javascript web-component custom-element native-web-component
我已阅读了类似问题的一些问题以及“帮助1”页面。不幸的是,我被困住了。
一个可能的原因可能是代理引起的,但是这里没有这样的代理。另外,当我从Eclipse更新它时,我PC中的所有maven项目都已成功更新。所以我放弃了这种可能性。
我检查的另一件事是在本地存储库中查找codehaus,然后找到它(C:\ Users \ myUser.m2 \ repository \ org \ codehaus \ mojo)。
另一种尝试,我尝试在设置中添加pluginGroups / pluginGroup。该项目是一个非常简单的问候词,仅使用Spring Batch和Tasklet及其execute方法。没有公共的static void main方法。
我添加了Matteo建议后,打印了整个错误的Command.exe:
C:\temp\TaskletJavaConfig\spring-batch-helloworld>mvn compile exec:java -e
Picked up JAVA_TOOL_OPTIONS: -agentlib:jvmhook
Picked up _JAVA_OPTIONS: -Xrunjvmhook -Xbootclasspath/a:C:\PROGRA~2\HP\QUICKT~1\
bin\JAVA_S~1\classes;C:\PROGRA~2\HP\QUICKT~1\bin\JAVA_S~1\classes\jasmine.jar
[INFO] Error stacktraces are turned on.
[INFO] Scanning for projects...
Downloading: https://repo.maven.apache.org/maven2/org/codehaus/mojo/exec-maven-p
lugin/1.1/exec-maven-plugin-1.1.pom
[WARNING] Failed to retrieve plugin descriptor for org.codehaus.mojo:exec-maven-
plugin:1.1: Plugin org.codehaus.mojo:exec-maven-plugin:1.1 or one of its depende
ncies could not be resolved: Failed to read artifact descriptor for org.codehaus
.mojo:exec-maven-plugin:jar:1.1
Downloading: https://repo.maven.apache.org/maven2/org/apache/maven/plugins/maven
-deploy-plugin/2.7/maven-deploy-plugin-2.7.pom …Run Code Online (Sandbox Code Playgroud) 我想了解如何使用 WebFlux Restcontroller 从 React 生成 json 流。我的第一个尝试迫使我解析为文本而不是 json,但我感觉我在做一些奇怪的事情。我决定投入精力阅读周围的示例,我发现一个返回 org.reactivestreams.Publisher 而不是返回 WebFlux 的示例。
我发现一个有点旧的主题可以帮助我(Unable to Consume Webflux Streaming Response in React-Native client),但它没有一个答案。好吧,它促使我在某个地方读到 React 可能不符合 Reactive,但这是一篇非常旧的文章。
从互联网上获取样本,如果你查看https://developer.okta.com/blog/2018/09/25/spring-webflux-websockets-react基本上你会发现:
WebFlux 生成 Json 但不流式传输:
RestController producing MediaType.APPLICATION_JSON_VALUE and returning org.reactivestreams.Publisher<the relevant pojo>
Run Code Online (Sandbox Code Playgroud)
React 消费 json:
async componentDidMount() {
const response = await fetch('relevant url');
const data = await response.json();
}
Run Code Online (Sandbox Code Playgroud)
我尝试过类似的方法,但有两个显着的区别:
我返回 MediaType.TEXT_EVENT_STREAM_VALUE 因为我相信我应该更喜欢返回流,因为我正在使用无阻塞代码,并且我知道返回 Webflux 更有意义,以便正确利用 Spring 5,而不是返回 org.reactivestreams.Publisher +。
我的 Webflux 休息控制器:
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.http.MediaType;
import org.springframework.web.bind.annotation.CrossOrigin;
import …Run Code Online (Sandbox Code Playgroud) reactive-programming async-await reactjs jsonstream spring-webflux
我有一个非常简单的节点服务,它公开了一个旨在使用服务器发送事件 (SSE) 连接的端点和一个非常基本的 ReactJs 客户端,通过 EventSource.onmessage 使用它。
首先,当我在 updateAmountState (Chrome Dev) 中设置调试点时,我看不到它被唤起。
其次,我得到 net::ERR_INCOMPLETE_CHUNKED_ENCODING 200(好的)。根据https://github.com/aspnet/KestrelHttpServer/issues/1858 “Chrome 中的 ERR_INCOMPLETE_CHUNKED_ENCODING 通常意味着在写入响应正文的过程中应用程序抛出了未捕获的异常”。然后我检查了服务器端,看看我是否发现了任何错误。好吧,我在 setTimeout(() => {... 的 server.js 中的几个地方设置了断点,我看到它定期运行。我希望每行只运行一次。所以看起来前端是尝试永久调用后端并收到一些错误。
整个应用程序,包括 ReactJs 的前端和 NodeJs 的服务器都可以在https://github.com/jimisdrpc/hello-pocker-coins 中找到。
后端:
const http = require("http");
http
.createServer((request, response) => {
console.log("Requested url: " + request.url);
if (request.url.toLowerCase() === "/coins") {
response.writeHead(200, {
Connection: "keep-alive",
"Content-Type": "text/event-stream",
"Cache-Control": "no-cache"
});
setTimeout(() => {
response.write('data: {"player": "Player1", "amount": "90"}');
response.write("\n\n");
}, 3000);
setTimeout(() => {
response.write('data: {"player": "Player2", …Run Code Online (Sandbox Code Playgroud) node.js google-chrome-devtools server-sent-events reactjs eventsource
我有一个挑战我的场景:Kafka 主题,我必须使用这些消息并通过 SSE 公开给 Web 组件。我已经在所有层中进行了几次尝试,以寻找一种稳定且更可靠的方法,或者至少有一些我可以更舒适地支持。
现在我创建了一个非常简单的 Kafka 主题并创建了两个不同的使用者,两者的工作方式似乎非常相似。一种是使用 org.springframework.kafka.annotation.KafkaListener 和其他 reactor.kafka.receiver.KafkaReceiver。
最终目标是通过 Spring WebFlux 公开一个事件,当消费者来自主题时发布消息。
如果我没记错的话,我在某处读到 Spring..KafkaListener 正在阻止代码,但据我所知,它不是。它只是一个与 Reactor..KafkaReceiver 完全相同的监听器。因为我正在编写一个非阻塞代码,所以我应该避免阻塞代码,但我不能在任何地方使用 Spring..KafkaListener 阻塞。
以下是产生完全相同结果的基本比较:
反应堆KafkaReceiver:
ReceiverOptions<Object, Object> consumerOptions = ReceiverOptions.create(consumerProps)
.subscription(Collections.singleton("test"))
.addAssignListener(partitions -> logger.debug("onPartitionsAssigned {}", partitions))
.addRevokeListener(partitions -> logger.debug("onPartitionsRevoked {}", partitions));
kafkaReceiver = KafkaReceiver.create(consumerOptions);
((Flux<ReceiverRecord>) kafkaReceiver.receive()).doOnNext(r -> {
logger.info(String.format("Consumed Message using KafkaListener -> %s", r.value()));
r.receiverOffset().acknowledge();
}).subscribe();
Run Code Online (Sandbox Code Playgroud)
春季卡夫卡听众:
@KafkaListener(topics = "test")
public void consume(String message) {
logger.info(String.format("Consumed Message using KafkaListener -> %s", message));
}
Run Code Online (Sandbox Code Playgroud)
如果要重现比较:
1 - clone https://github.com/jimisdrpc/simplest-comparison-kafkaconsumer
2 …Run Code Online (Sandbox Code Playgroud) web-component server-sent-events apache-kafka project-reactor spring-webflux
我正在尝试一个非常简单的教程,解释如何将 docker-compose 转换为 minishift (Minishift 和 Kompose。我尝试转换并推送 docker-compose.yml 示例
\n\nversion: "2"\n\nservices:\n\n redis-master:\n image: k8s.gcr.io/redis:e2e \n ports:\n - "6379"\n\n redis-slave:\n image: gcr.io/google_samples/gb-redisslave:v1\n ports:\n - "6379"\n environment:\n - GET_HOSTS_FROM=dns\n\n frontend:\n image: gcr.io/google-samples/gb-frontend:v4\n ports:\n - "80:80"\n environment:\n - GET_HOSTS_FROM=dns\n labels:\n kompose.service.type: LoadBalancer\nRun Code Online (Sandbox Code Playgroud)\n\n从这些日志中可以看到,我成功地编写并推送了:
\n\nC:\\Users\\Cast\\docker-compose-to-minishift>kompose-windows-amd64 up --provider=openshift\n[36mINFO[0m We are going to create OpenShift DeploymentConfigs, Services and PersistentVolumeClaims for your Dockerized application.\nIf you need different kind of resources, use the \'kompose convert\' and \'oc create -f\' commands instead.\n\n[36mINFO[0m Deploying application …Run Code Online (Sandbox Code Playgroud) 直接的问题是:为什么 Gradle 没有解决我添加的这个依赖项
dependencies {
//kafka-protobuf-serializer
implementation("io.confluent:kafka-protobuf-serializer:6.0.0")
}
Run Code Online (Sandbox Code Playgroud)
?
根据mvn,这是我在 build.gradle 中添加此类依赖项的方式
compile group: 'io.confluent', name: 'kafka-protobuf-serializer', version: '6.0.0'
Run Code Online (Sandbox Code Playgroud)
我所有的其他依赖项都添加了“实现...”。到现在为止还挺好。但是对于这个特殊的我得到了
Execution failed for task ':extractIncludeProto'.
> Could not resolve all files for configuration ':compileProtoPath'.
> Could not find io.confluent:kafka-protobuf-serializer:6.0.0.
Searched in the following locations:
- file:/C:/Users/Cast/.m2/repository/io/confluent/kafka-protobuf-serializer/6.0.0/kafka-protobuf-serializer-6.0.0.pom
- https://jcenter.bintray.com/io/confluent/kafka-protobuf-serializer/6.0.0/kafka-protobuf-serializer-6.0.0.pom
Required by:
project :
Possible solution:
- Declare repository providing the artifact, see the documentation at https://docs.gradle.org/current/userguide/declaring_repositories.html
Run Code Online (Sandbox Code Playgroud)
我在这里缺少什么或弄乱了什么?
这是整个 build.gradle
plugins {
id "org.jetbrains.kotlin.jvm" version "1.3.72"
id "org.jetbrains.kotlin.kapt" version "1.3.72" …Run Code Online (Sandbox Code Playgroud) 目的:我想编码一个端点,它利用轻量协程消耗另一个端点,假设我正在编码一个轻量异步端点客户端。
我的背景:第一次尝试使用 Kotlin Coroutine。我最近几天研究并四处寻找。我发现很多文章解释如何在 Android 中使用 Coroutine,但很少有其他文章解释如何在主函数中使用 Coroutine。不幸的是,我没有找到解释如何使用协程编写控制器端点的文章,如果我正在做一些不推荐的事情,它就会在我的脑海中敲响警钟。
目前情况:我使用协程成功创建了几种方法,但我想知道哪种方法最适合传统的 GET。最重要的是,我想知道如何正确处理异常。
主要问题:推荐以下方法中的哪一种以及我应该关注哪种异常处理?
相关的第二个问题:有什么区别
fun someMethodWithRunBlocking(): String? = runBlocking {
return@runBlocking ...
}
Run Code Online (Sandbox Code Playgroud)
和
suspend fun someMethodWithSuspendModifier(): String?{
return ...
}
Run Code Online (Sandbox Code Playgroud)
下面的所有尝试都在工作并返回 json 响应,但我不知道端点方法上的“runBlocking”和返回“return@runBlocking”是否会给我带来一些负面的缺点。
控制器(端点)
package com.tolearn.controller
import com.tolearn.service.DemoService
import io.micronaut.http.MediaType
import io.micronaut.http.annotation.Controller
import io.micronaut.http.annotation.Get
import io.micronaut.http.annotation.Produces
import kotlinx.coroutines.Deferred
import kotlinx.coroutines.*
import kotlinx.coroutines.runBlocking
import java.net.http.HttpResponse
import javax.inject.Inject
@Controller("/tolearn")
class DemoController {
@Inject
lateinit var demoService: DemoService
//APPROACH 1:
//EndPoint method with runBlocking CoroutineScope
//Using Deferred.await
//Using return@runBlocking
@Get("/test1")
@Produces(MediaType.TEXT_PLAIN)
fun getWithRunBlockingAndDeferred(): String? …Run Code Online (Sandbox Code Playgroud) 根据/spring-kafka/docs/2.4.4.RELEASE/,关于Kafka的新特性旨在否定确认,现在由Spring-Kafka支持
"... 从 2.3 版本开始,Acknowledgment 接口有两个额外的方法 nack(long sleep) 和 nack(int index, long sleep)。第一个用于记录侦听器,第二个用于批处理侦听器。调用您的侦听器类型的错误方法将引发 IllegalStateException。
...
使用记录侦听器,当调用 nack() 时,将提交任何挂起的偏移量,丢弃上次轮询的剩余记录,并在其分区上执行查找,以便在下一次轮询时重新发送失败的记录和未处理的记录( )。通过设置 sleep 参数,可以在重新传递之前暂停使用者线程。这与在容器配置有 SeekToCurrentErrorHandler 时抛出异常的功能类似。”
好吧,如果消费者方面发生了一些错误,比如说无法保存在数据库上,假设消费者没有acknowledge.acknowledge(),据我所知,消息仍在轮询中,它将再次被读取/使用。我想有人可以说,使用 nack(..., some time) 消费者可以睡觉,让有机会稍后再次阅读/消费并且不会遇到错误。如果继续听这个话题不是问题,我的直接问题是:
使用 nack 而不是简单地不承认还有什么意义吗?
据我所知,无论如何,该消息将在池中保留的时间长于 nack sleep 的时间。因此,顺便说一下,如果消费者不断尝试获取消息并保存消息,假设问题在睡眠时间内得到解决,它就会成功。
一个周围的点或优势是,以某种方式生产者得到通知使用 nack。如果是这样,我可以在某些特定场景中找到一些价值。假设使用 Log Compation(仅对最后一条消息状态感兴趣)或 Kafka 作为长期存储服务(我猜未来版本将提供此服务 - KIP 405)
关于更一般的异常,我倾向于遵循配置 SeekToCurrentErrorHandler 并抛出异常的方法
java spring-integration apache-kafka kafka-consumer-api spring-kafka
apache-kafka ×3
java ×2
java-11 ×2
maven ×2
reactjs ×2
apache2 ×1
async-await ×1
docker ×1
eventsource ×1
gradle ×1
html ×1
javascript ×1
jsonstream ×1
kotlin ×1
kubernetes ×1
maven-3 ×1
node.js ×1
openshift ×1
spring ×1
spring-kafka ×1