mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-21 05:44:03 +08:00
* Add namespace v2 * Add NamespaceRpcHook * Refector extends header * Use Boolean in request header to remove unnecessary encode * Add unit test * Add NamespaceRpcHookTest * Remove GrpcConverter#wrapResourceWithNamespace * Optimize readability of RpcRequestHeader
This commit is contained in:
@@ -141,7 +141,6 @@ public class TransactionProducer {
|
||||
}
|
||||
final TransactionListener transactionCheckListener = new TransactionListenerImpl(statsBenchmark, config);
|
||||
final TransactionMQProducer producer = new TransactionMQProducer(
|
||||
null,
|
||||
"benchmark_transaction_producer",
|
||||
rpcHook,
|
||||
config.msgTraceEnable,
|
||||
|
||||
+2
-1
@@ -33,7 +33,8 @@ public class ProducerWithNamespace {
|
||||
|
||||
public static void main(String[] args) throws Exception {
|
||||
|
||||
DefaultMQProducer producer = new DefaultMQProducer(NAMESPACE, PRODUCER_GROUP);
|
||||
DefaultMQProducer producer = new DefaultMQProducer(PRODUCER_GROUP);
|
||||
producer.setNamespaceV2(NAMESPACE);
|
||||
|
||||
producer.setNamesrvAddr(DEFAULT_NAMESRVADDR);
|
||||
producer.start();
|
||||
|
||||
+2
-1
@@ -34,7 +34,8 @@ public class PullConsumerWithNamespace {
|
||||
private static final Map<MessageQueue, Long> OFFSET_TABLE = new HashMap<>();
|
||||
|
||||
public static void main(String[] args) throws Exception {
|
||||
DefaultMQPullConsumer pullConsumer = new DefaultMQPullConsumer(NAMESPACE, CONSUMER_GROUP);
|
||||
DefaultMQPullConsumer pullConsumer = new DefaultMQPullConsumer(CONSUMER_GROUP);
|
||||
pullConsumer.setNamespaceV2(NAMESPACE);
|
||||
pullConsumer.setNamesrvAddr(DEFAULT_NAMESRVADDR);
|
||||
pullConsumer.start();
|
||||
|
||||
|
||||
+2
-1
@@ -27,7 +27,8 @@ public class PushConsumerWithNamespace {
|
||||
public static final String TOPIC = "NAMESPACE_TOPIC";
|
||||
|
||||
public static void main(String[] args) throws Exception {
|
||||
DefaultMQPushConsumer defaultMQPushConsumer = new DefaultMQPushConsumer(NAMESPACE, CONSUMER_GROUP);
|
||||
DefaultMQPushConsumer defaultMQPushConsumer = new DefaultMQPushConsumer(CONSUMER_GROUP);
|
||||
defaultMQPushConsumer.setNamespaceV2(NAMESPACE);
|
||||
defaultMQPushConsumer.setNamesrvAddr(DEFAULT_NAMESRVADDR);
|
||||
defaultMQPushConsumer.subscribe(TOPIC, "*");
|
||||
defaultMQPushConsumer.registerMessageListener((MessageListenerConcurrently) (msgs, context) -> {
|
||||
|
||||
@@ -34,7 +34,7 @@ public class TraceProducer {
|
||||
|
||||
public static void main(String[] args) throws MQClientException, InterruptedException {
|
||||
|
||||
DefaultMQProducer producer = new DefaultMQProducer(PRODUCER_GROUP, true);
|
||||
DefaultMQProducer producer = new DefaultMQProducer(PRODUCER_GROUP, true, null);
|
||||
|
||||
// Uncomment the following line while debugging, namesrvAddr should be set to your local address
|
||||
// producer.setNamesrvAddr(DEFAULT_NAMESRVADDR);
|
||||
|
||||
+1
-1
@@ -31,7 +31,7 @@ public class TracePushConsumer {
|
||||
|
||||
public static void main(String[] args) throws InterruptedException, MQClientException {
|
||||
// Here,we use the default message track trace topic name
|
||||
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer(CONSUMER_GROUP, true);
|
||||
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer(CONSUMER_GROUP, true, null);
|
||||
|
||||
// Uncomment the following line while debugging, namesrvAddr should be set to your local address
|
||||
// consumer.setNamesrvAddr(DEFAULT_NAMESRVADDR);
|
||||
|
||||
Reference in New Issue
Block a user