rocketmq-5.3.2集群A
IP: X.X.X.A (/rocketmq/bin/dledger/fast-try.sh)
rocketmq-5.3.2集群B
IP: X.X.X.B (/rocketmq/bin/dledger/fast-try.sh)
rocketmq-connect:
IP: X.X.X.A
单机模式启动
sh bin/connect-standalone.sh -c conf/connect-standalone.conf &
rocketmq-replicator:
ps:参数说明只有四个必填,但是实际启动region,cloud不填会报错500
curl -X POST -H "Content-Type: application/json" http://127.0.0.1:8082/connectors/full_msg_replicator -d '{
"connector.class": "org.apache.rocketmq.replicator.ReplicatorSourceConnector",
"src.endpoint": "X.X.X.A:9876",
"src.cluster": "RaftCluster",
"src.region": "regionA",
"src.cloud": "cloudA",
"src.topictags": "TopicTest,*",
"src.acl.enable": "false",
"max.task": "4",
"dest.endpoint": "X.X.X.B:9876",
"dest.cluster": "RaftCluster",
"dest.region": "regionB",
"dest.cloud": "cloudB",
"dest.topic": "TopicTest",
"dest.acl.enable": "false",
"errors.tolerance": "none"
}'
查看集群A消费组,CID_RMQ_SYS_REPLICATOR_full_msg_replicator开始正常消费,但是集群B的TopicTest的消息数并未增加,且集群A的TopicTest消息数开始增加。
查看connect运行时的日志,其中对应发送消息到的就是集群A,明明已经配置了dest.endpoint到了集群B的路由地址,为什么会出现这种问题?应该如何修改?
org.apache.rocketmq.client.exception.MQBrokerException: CODE: 2 DESC: [PC_SYNCHRONIZED]broker busy, start flow control for a while
BROKER: X.X.X.A:30911
For more information, please visit the url, http://rocketmq.apache.org/docs/faq/
at org.apache.rocketmq.client.impl.MQClientAPIImpl.processSendResponse(MQClientAPIImpl.java:675)
at org.apache.rocketmq.client.impl.MQClientAPIImpl.access$000(MQClientAPIImpl.java:175)
at org.apache.rocketmq.client.impl.MQClientAPIImpl$1.operationComplete(MQClientAPIImpl.java:564)
at org.apache.rocketmq.remoting.netty.ResponseFuture.executeInvokeCallback(ResponseFuture.java:54)
at org.apache.rocketmq.remoting.netty.NettyRemotingAbstract$2.run(NettyRemotingAbstract.java:321)
at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:515)
at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264)
at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128)
at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628)
at java.base/java.lang.Thread.run(Thread.java:829)
rocketmq-5.3.2集群A
IP: X.X.X.A (/rocketmq/bin/dledger/fast-try.sh)
rocketmq-5.3.2集群B
IP: X.X.X.B (/rocketmq/bin/dledger/fast-try.sh)
rocketmq-connect:
IP: X.X.X.A
单机模式启动
sh bin/connect-standalone.sh -c conf/connect-standalone.conf &
rocketmq-replicator:
ps:参数说明只有四个必填,但是实际启动region,cloud不填会报错500
curl -X POST -H "Content-Type: application/json" http://127.0.0.1:8082/connectors/full_msg_replicator -d '{
"connector.class": "org.apache.rocketmq.replicator.ReplicatorSourceConnector",
"src.endpoint": "X.X.X.A:9876",
"src.cluster": "RaftCluster",
"src.region": "regionA",
"src.cloud": "cloudA",
"src.topictags": "TopicTest,*",
"src.acl.enable": "false",
"max.task": "4",
"dest.endpoint": "X.X.X.B:9876",
"dest.cluster": "RaftCluster",
"dest.region": "regionB",
"dest.cloud": "cloudB",
"dest.topic": "TopicTest",
"dest.acl.enable": "false",
"errors.tolerance": "none"
}'
查看集群A消费组,CID_RMQ_SYS_REPLICATOR_full_msg_replicator开始正常消费,但是集群B的TopicTest的消息数并未增加,且集群A的TopicTest消息数开始增加。
查看connect运行时的日志,其中对应发送消息到的就是集群A,明明已经配置了dest.endpoint到了集群B的路由地址,为什么会出现这种问题?应该如何修改?
org.apache.rocketmq.client.exception.MQBrokerException: CODE: 2 DESC: [PC_SYNCHRONIZED]broker busy, start flow control for a while
BROKER: X.X.X.A:30911
For more information, please visit the url, http://rocketmq.apache.org/docs/faq/
at org.apache.rocketmq.client.impl.MQClientAPIImpl.processSendResponse(MQClientAPIImpl.java:675)
at org.apache.rocketmq.client.impl.MQClientAPIImpl.access$000(MQClientAPIImpl.java:175)
at org.apache.rocketmq.client.impl.MQClientAPIImpl$1.operationComplete(MQClientAPIImpl.java:564)
at org.apache.rocketmq.remoting.netty.ResponseFuture.executeInvokeCallback(ResponseFuture.java:54)
at org.apache.rocketmq.remoting.netty.NettyRemotingAbstract$2.run(NettyRemotingAbstract.java:321)
at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:515)
at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264)
at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128)
at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628)
at java.base/java.lang.Thread.run(Thread.java:829)