Skip to content
This repository was archived by the owner on Jan 24, 2024. It is now read-only.
This repository was archived by the owner on Jan 24, 2024. It is now read-only.

[BUG] Kafka producer client can not connect to kop (Removing node xxxx:9092 (id: 2057312963 rack: null) from least loaded node selection since it is neither ready for sending or connecting)[BUG] #618

Description

@xiaotongwang1
"pulsar-io-26-32" #287 prio=5 os_prio=0 tid=0x00007f10e8767000 nid=0x4ae waiting on condition [0x00007f10c95ba000]
   java.lang.Thread.State: WAITING (parking)
	at sun.misc.Unsafe.park(Native Method)
	- parking to wait for  <0x000000052b910ae0> (a org.apache.pulsar.common.util.collections.ConcurrentLongHashMap$Section)
	at java.util.concurrent.locks.StampedLock.acquireWrite(StampedLock.java:1119)
	at java.util.concurrent.locks.StampedLock.writeLock(StampedLock.java:354)
	at org.apache.pulsar.common.util.collections.ConcurrentLongHashMap$Section.remove(ConcurrentLongHashMap.java:340)
	at org.apache.pulsar.common.util.collections.ConcurrentLongHashMap$Section.access$300(ConcurrentLongHashMap.java:194)
	at org.apache.pulsar.common.util.collections.ConcurrentLongHashMap.remove(ConcurrentLongHashMap.java:141)
	at io.streamnative.pulsar.handlers.kop.KafkaTopicConsumerManager.deleteOneExpiredCursor(KafkaTopicConsumerManager.java:91)
	at io.streamnative.pulsar.handlers.kop.KafkaTopicConsumerManager.lambda$deleteExpiredCursor$0(KafkaTopicConsumerManager.java:76)
	at io.streamnative.pulsar.handlers.kop.KafkaTopicConsumerManager$$Lambda$1120/2085128569.accept(Unknown Source)
	at org.apache.pulsar.common.util.collections.ConcurrentLongHashMap$Section.forEach(ConcurrentLongHashMap.java:431)
	at org.apache.pulsar.common.util.collections.ConcurrentLongHashMap.forEach(ConcurrentLongHashMap.java:164)
	at io.streamnative.pulsar.handlers.kop.KafkaTopicConsumerManager.deleteExpiredCursor(KafkaTopicConsumerManager.java:74)
	at io.streamnative.pulsar.handlers.kop.KafkaTopicManager.lambda$null$0(KafkaTopicManager.java:100)
	at io.streamnative.pulsar.handlers.kop.KafkaTopicManager$$Lambda$1119/1991303689.accept(Unknown Source)
	at java.util.concurrent.ConcurrentHashMap$ValuesView.forEach(ConcurrentHashMap.java:4707)
	at io.streamnative.pulsar.handlers.kop.KafkaTopicManager.lambda$new$1(KafkaTopicManager.java:98)
	at io.streamnative.pulsar.handlers.kop.KafkaTopicManager$$Lambda$893/144794386.run(Unknown Source)
	at io.netty.util.concurrent.PromiseTask.runTask(PromiseTask.java:98)
	at io.netty.util.concurrent.ScheduledFutureTask.run(ScheduledFutureTask.java:176)
	at io.netty.util.concurrent.AbstractEventExecutor.safeExecute(AbstractEventExecutor.java:164)
	at io.netty.util.concurrent.SingleThreadEventExecutor.runAllTasks(SingleThreadEventExecutor.java:472)
	at io.netty.channel.epoll.EpollEventLoop.run(EpollEventLoop.java:384)
	at io.netty.util.concurrent.SingleThreadEventExecutor$4.run(SingleThreadEventExecutor.java:989)
	at io.netty.util.internal.ThreadExecutorMap$2.run(ThreadExecutorMap.java:74)
	at io.netty.util.concurrent.FastThreadLocalRunnable.run(FastThreadLocalRunnable.java:30)
	at java.lang.Thread.run(Thread.java:748)


"pulsar-io-26-31" #286 prio=5 os_prio=0 tid=0x00007f10e8445000 nid=0x4ad waiting on condition [0x00007f10c96bb000]
   java.lang.Thread.State: WAITING (parking)
	at sun.misc.Unsafe.park(Native Method)
	- parking to wait for  <0x000000052b910ae0> (a org.apache.pulsar.common.util.collections.ConcurrentLongHashMap$Section)
	at java.util.concurrent.locks.StampedLock.acquireWrite(StampedLock.java:1119)
	at java.util.concurrent.locks.StampedLock.writeLock(StampedLock.java:354)
	at org.apache.pulsar.common.util.collections.ConcurrentLongHashMap$Section.remove(ConcurrentLongHashMap.java:340)
	at org.apache.pulsar.common.util.collections.ConcurrentLongHashMap$Section.access$300(ConcurrentLongHashMap.java:194)
	at org.apache.pulsar.common.util.collections.ConcurrentLongHashMap.remove(ConcurrentLongHashMap.java:141)
	at io.streamnative.pulsar.handlers.kop.KafkaTopicConsumerManager.remove(KafkaTopicConsumerManager.java:142)
	at io.streamnative.pulsar.handlers.kop.MessageFetchContext.lambda$null$2(MessageFetchContext.java:146)
	at io.streamnative.pulsar.handlers.kop.MessageFetchContext$$Lambda$945/1999044484.apply(Unknown Source)
	at java.util.stream.ReferencePipeline$3$1.accept(ReferencePipeline.java:193)
	at java.util.HashMap$EntrySpliterator.forEachRemaining(HashMap.java:1699)
	at java.util.stream.AbstractPipeline.copyInto(AbstractPipeline.java:482)
	at java.util.stream.AbstractPipeline.wrapAndCopyInto(AbstractPipeline.java:472)
	at java.util.stream.ReduceOps$ReduceOp.evaluateSequential(ReduceOps.java:708)
	at java.util.stream.AbstractPipeline.evaluate(AbstractPipeline.java:234)
	at java.util.stream.ReferencePipeline.collect(ReferencePipeline.java:566)
	at io.streamnative.pulsar.handlers.kop.MessageFetchContext.lambda$handleFetch$4(MessageFetchContext.java:167)
	at io.streamnative.pulsar.handlers.kop.MessageFetchContext$$Lambda$944/1214167246.accept(Unknown Source)
	at java.util.concurrent.CompletableFuture.uniWhenComplete(CompletableFuture.java:774)
	at java.util.concurrent.CompletableFuture.uniWhenCompleteStage(CompletableFuture.java:792)
	at java.util.concurrent.CompletableFuture.whenComplete(CompletableFuture.java:2153)
	at io.streamnative.pulsar.handlers.kop.MessageFetchContext.handleFetch(MessageFetchContext.java:109)
	at io.streamnative.pulsar.handlers.kop.KafkaRequestHandler.handleFetchRequest(KafkaRequestHandler.java:966)
	at io.streamnative.pulsar.handlers.kop.KafkaCommandDecoder.channelRead(KafkaCommandDecoder.java:203)
	at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:379)
	at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:365)
	at io.netty.channel.AbstractChannelHandlerContext.fireChannelRead(AbstractChannelHandlerContext.java:357)
	at io.netty.handler.codec.ByteToMessageDecoder.fireChannelRead(ByteToMessageDecoder.java:324)
	at io.netty.handler.codec.ByteToMessageDecoder.channelRead(ByteToMessageDecoder.java:296)
	at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:379)
	at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:365)
	at io.netty.channel.AbstractChannelHandlerContext.fireChannelRead(AbstractChannelHandlerContext.java:357)
	at io.netty.channel.DefaultChannelPipeline$HeadContext.channelRead(DefaultChannelPipeline.java:1410)
	at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:379)
	at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:365)
	at io.netty.channel.DefaultChannelPipeline.fireChannelRead(DefaultChannelPipeline.java:919)
	at io.netty.channel.epoll.AbstractEpollStreamChannel$EpollStreamUnsafe.epollInReady(AbstractEpollStreamChannel.java:792)
	at io.netty.channel.epoll.EpollEventLoop.processReady(EpollEventLoop.java:475)
	at io.netty.channel.epoll.EpollEventLoop.run(EpollEventLoop.java:378)
	at io.netty.util.concurrent.SingleThreadEventExecutor$4.run(SingleThreadEventExecutor.java:989)
	at io.netty.util.internal.ThreadExecutorMap$2.run(ThreadExecutorMap.java:74)
	at io.netty.util.concurrent.FastThreadLocalRunnable.run(FastThreadLocalRunnable.java:30)
	at java.lang.Thread.run(Thread.java:748)

"pulsar-io-26-30" #285 prio=5 os_prio=0 tid=0x00007f10e8443800 nid=0x4ac waiting on condition [0x00007f10c97bd000]
   java.lang.Thread.State: WAITING (parking)
	at sun.misc.Unsafe.park(Native Method)
	- parking to wait for  <0x000000052b9108a0> (a java.util.concurrent.locks.ReentrantReadWriteLock$NonfairSync)
	at java.util.concurrent.locks.LockSupport.park(LockSupport.java:175)
	at java.util.concurrent.locks.AbstractQueuedSynchronizer.parkAndCheckInterrupt(AbstractQueuedSynchronizer.java:838)
	at java.util.concurrent.locks.AbstractQueuedSynchronizer.acquireQueued(AbstractQueuedSynchronizer.java:872)
	at java.util.concurrent.locks.AbstractQueuedSynchronizer.acquire(AbstractQueuedSynchronizer.java:1201)
	at java.util.concurrent.locks.ReentrantReadWriteLock$WriteLock.lock(ReentrantReadWriteLock.java:943)
	at io.streamnative.pulsar.handlers.kop.KafkaTopicConsumerManager.close(KafkaTopicConsumerManager.java:244)
	at io.streamnative.pulsar.handlers.kop.KafkaTopicManager.close(KafkaTopicManager.java:336)
	- locked <0x000000052e940220> (a io.streamnative.pulsar.handlers.kop.KafkaTopicManager)
	at io.streamnative.pulsar.handlers.kop.KafkaRequestHandler.close(KafkaRequestHandler.java:212)
	at io.streamnative.pulsar.handlers.kop.KafkaRequestHandler.channelInactive(KafkaRequestHandler.java:202)
	at io.netty.channel.AbstractChannelHandlerContext.invokeChannelInactive(AbstractChannelHandlerContext.java:262)
	at io.netty.channel.AbstractChannelHandlerContext.invokeChannelInactive(AbstractChannelHandlerContext.java:248)
	at io.netty.channel.AbstractChannelHandlerContext.fireChannelInactive(AbstractChannelHandlerContext.java:241)
	at io.netty.handler.codec.ByteToMessageDecoder.channelInputClosed(ByteToMessageDecoder.java:389)
	at io.netty.handler.codec.ByteToMessageDecoder.channelInactive(ByteToMessageDecoder.java:354)
	at io.netty.channel.AbstractChannelHandlerContext.invokeChannelInactive(AbstractChannelHandlerContext.java:262)
	at io.netty.channel.AbstractChannelHandlerContext.invokeChannelInactive(AbstractChannelHandlerContext.java:248)
	at io.netty.channel.AbstractChannelHandlerContext.fireChannelInactive(AbstractChannelHandlerContext.java:241)
	at io.netty.channel.DefaultChannelPipeline$HeadContext.channelInactive(DefaultChannelPipeline.java:1405)
	at io.netty.channel.AbstractChannelHandlerContext.invokeChannelInactive(AbstractChannelHandlerContext.java:262)
	at io.netty.channel.AbstractChannelHandlerContext.invokeChannelInactive(AbstractChannelHandlerContext.java:248)
	at io.netty.channel.DefaultChannelPipeline.fireChannelInactive(DefaultChannelPipeline.java:901)
	at io.netty.channel.AbstractChannel$AbstractUnsafe$8.run(AbstractChannel.java:818)
	at io.netty.util.concurrent.AbstractEventExecutor.safeExecute(AbstractEventExecutor.java:164)
	at io.netty.util.concurrent.SingleThreadEventExecutor.runAllTasks(SingleThreadEventExecutor.java:472)
	at io.netty.channel.epoll.EpollEventLoop.run(EpollEventLoop.java:384)
	at io.netty.util.concurrent.SingleThreadEventExecutor$4.run(SingleThreadEventExecutor.java:989)
	at io.netty.util.internal.ThreadExecutorMap$2.run(ThreadExecutorMap.java:74)
	at io.netty.util.concurrent.FastThreadLocalRunnable.run(FastThreadLocalRunnable.java:30)
	at java.lang.Thread.run(Thread.java:748)

Metadata

Metadata

Assignees

No one assigned

    Labels

    Type

    No type

    Projects

    No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions