apache / apache/seatunnel

[Bug] [RocketMQ Source] Only one parallism comsumer data in Fink Engine

Open
#9,706 3 comments 0 reactions 1 assignee Claimed by @zhangshenghang View on GitHub
bug
Dominant language
Java
Stars
9.7k
Forks
2.4k
Avg merge
3d 9h
Merged PRs (30d)
204

Description

### Search before asking

- [x] I had searched in the [issues](https://github.com/apache/seatunnel/issues?q=is%3Aissue+label%3A%22bug%22) and found no similar issues.

### What happened

Using RocketMQ Source consumer data,flink backprssure always very high,and only subtask0 consumer the data Image

### SeaTunnel Version

2.3.11

### SeaTunnel Config

```conf
{
"env": {
"job.mode": "streaming",
"checkpoint.interval": "8000",
"parallelism": 8,
"checkpoint.timeout": 600000,
"flink.jobmanager.memory.process.size": "1024m",
"flink.taskmanager.memory.process.size": "2048m",
"flink.taskmanager.memory.managed.fraction": "0.05",
"flink.taskmanager.numberOfTaskSlots": 4
},
"source": [
{
"plugin_name": "Rocketmq",
"plugin_output": "source_tabl1",
"name.srv.addr": "localhost:9876",
"start.mode": "CONSUME_FROM_LAST_OFFSET",
"consumer.group": "tpws-seatunnel-source-record-topic-consumer",
"topics": "tpws-op-new-record-topic",
"schema": {
"fields": {
"merged": "map"
}
}
}
],
"sink": [
{
"plugin_name": "Doris",
"plugin_input": "source_tabl1",
"doris.config": {
"format": "json",
"read_json_by_line": "true"
},
"fenodes": "localhost:8030",
"password": "********",
"username": "root",
"database": "tpws",
"table": "track_tws_detail",
"sink.label-prefix": "track_tws_detail",
"sink.enable-2pc": "true",
"sink.enable-delete": "true",
"sink.check-interval": 3000
}
]
}
```

### Running Command

```shell
xxx
```

### Error Exception

```log
java.lang.InterruptedException: null
at java.util.concurrent.locks.AbstractQueuedSynchronizer.doAcquireSharedNanos(AbstractQueuedSynchronizer.java:1039) ~[?:1.8.0_442]
at java.util.concurrent.locks.AbstractQueuedSynchronizer.tryAcquireSharedNanos(AbstractQueuedSynchronizer.java:1328) ~[?:1.8.0_442]
at java.util.concurrent.CountDownLatch.await(CountDownLatch.java:277) ~[?:1.8.0_442]
at org.apache.rocketmq.remoting.netty.ResponseFuture.waitResponse(ResponseFuture.java:71) ~[blob_p-c8e77e77a79f5fe5114e98aacb4db8fb43cdc369-92b33b3625a480b0051520f9adff7119:2.3.12-SNAPSHOT]
at org.apache.rocketmq.remoting.netty.NettyRemotingAbstract.invokeSyncImpl(NettyRemotingAbstract.java:435) ~[blob_p-c8e77e77a79f5fe5114e98aacb4db8fb43cdc369-92b33b3625a480b0051520f9adff7119:2.3.12-SNAPSHOT]
at org.apache.rocketmq.remoting.netty.NettyRemotingClient.invokeSync(NettyRemotingClient.java:390) ~[blob_p-c8e77e77a79f5fe5114e98aacb4db8fb43cdc369-92b33b3625a480b0051520f9adff7119:2.3.12-SNAPSHOT]
at org.apache.rocketmq.client.impl.MQClientAPIImpl.pullMessageSync(MQClientAPIImpl.java:779) ~[blob_p-c8e77e77a79f5fe5114e98aacb4db8fb43cdc369-92b33b3625a480b0051520f9adff7119:2.3.12-SNAPSHOT]
at org.apache.rocketmq.client.impl.MQClientAPIImpl.pullMessage(MQClientAPIImpl.java:733) ~[blob_p-c8e77e77a79f5fe5114e98aacb4db8fb43cdc369-92b33b3625a480b0051520f9adff7119:2.3.12-SNAPSHOT]
at org.apache.rocketmq.client.impl.consumer.PullAPIWrapper.pullKernelImpl(PullAPIWrapper.java:200) ~[blob_p-c8e77e77a79f5fe5114e98aacb4db8fb43cdc369-92b33b3625a480b0051520f9adff7119:2.3.12-SNAPSHOT]
at org.apache.rocketmq.client.impl.consumer.DefaultLitePullConsumerImpl.pullSyncImpl(DefaultLitePullConsumerImpl.java:930) ~[blob_p-c8e77e77a79f5fe5114e98aacb4db8fb43cdc369-92b33b3625a480b0051520f9adff7119:2.3.12-SNAPSHOT]
at org.apache.rocketmq.client.impl.consumer.DefaultLitePullConsumerImpl.pull(DefaultLitePullConsumerImpl.java:905) ~[blob_p-c8e77e77a79f5fe5114e98aacb4db8fb43cdc369-92b33b3625a480b0051520f9adff7119:2.3.12-SNAPSHOT]
at org.apache.rocketmq.client.impl.consumer.DefaultLitePullConsumerImpl.pull(DefaultLitePullConsumerImpl.java:900) ~[blob_p-c8e77e77a79f5fe5114e98aacb4db8fb43cdc369-92b33b3625a480b0051520f9adff7119:2.3.12-SNAPSHOT]
at org.apache.rocketmq.client.impl.consumer.DefaultLitePullConsumerImpl.access$1400(DefaultLitePullConsumerImpl.java:77) ~[blob_p-c8e77e77a79f5fe5114e98aacb4db8fb43cdc369-92b33b3625a480b0051520f9adff7119:2.3.12-SNAPSHOT]
at org.apache.rocketmq.client.impl.consumer.DefaultLitePullConsumerImpl$PullTaskImpl.run(DefaultLitePullConsumerImpl.java:849) [blob_p-c8e77e77a79f5fe5114e98aacb4db8fb43cdc369-92b33b3625a480b0051520f9adff7119:2.3.12-SNAPSHOT]
at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511) [?:1.8.0_442]
at java.util.concurrent.FutureTask.run(FutureTask.java:266) [?:1.8.0_442]
at java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.access$201(ScheduledThreadPoolExecutor.java:180) [?:1.8.0_442]
at java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.run(ScheduledThreadPoolExecutor.java:293) [?:1.8.0_442]
at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149) [?:1.8.0_442]
at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624) [?:1.8.0_442]
at java.lang.Thread.run(Thread.java:750) [?:1.8.0_442]
```

### Zeta or Flink or Spark Version

_No response_

### Java or Scala Version

_No response_

### Screenshots

_No response_

### Are you willing to submit PR?

- [x] Yes I am willing to submit a PR!

### Code of Conduct

- [x] I agree to follow this project's [Code of Conduct](https://www.apache.org/foundation/policies/conduct)

Contributor guide

No contributing guide indexed for this repository

Assessment

This issue has not been assessed yet.

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.