diff --git a/example/src/main/java/org/apache/rocketmq/example/batch/SimpleBatchProducer.java b/example/src/main/java/org/apache/rocketmq/example/batch/SimpleBatchProducer.java new file mode 100644 index 0000000000000000000000000000000000000000..a8609e793fadc0c3813cf573048fc87ccd06af6b --- /dev/null +++ b/example/src/main/java/org/apache/rocketmq/example/batch/SimpleBatchProducer.java @@ -0,0 +1,42 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.rocketmq.example.batch; + +import java.util.ArrayList; +import java.util.List; +import org.apache.rocketmq.client.producer.DefaultMQProducer; +import org.apache.rocketmq.common.message.Message; + +public class SimpleBatchProducer { + + + public static void main(String[] args) throws Exception { + DefaultMQProducer producer = new DefaultMQProducer("BatchProducerGroupName"); + producer.start(); + + //If you just send messages of no more than 1MiB at a time, it is easy to use batch + //Messages of the same batch should have: same topic, same waitStoreMsgOK and no schedule support + String topic = "BatchTest"; + List messages = new ArrayList<>(); + messages.add(new Message(topic, "Tag", "OrderID001", "Hello world 0".getBytes())); + messages.add(new Message(topic, "Tag", "OrderID002", "Hello world 1".getBytes())); + messages.add(new Message(topic, "Tag", "OrderID003", "Hello world 2".getBytes())); + + producer.send(messages); + } +} diff --git a/example/src/main/java/org/apache/rocketmq/example/batch/SplitBatchProducer.java b/example/src/main/java/org/apache/rocketmq/example/batch/SplitBatchProducer.java new file mode 100644 index 0000000000000000000000000000000000000000..8809a11512c5366e0bb76948e1375a65464d20e2 --- /dev/null +++ b/example/src/main/java/org/apache/rocketmq/example/batch/SplitBatchProducer.java @@ -0,0 +1,97 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.rocketmq.example.batch; + +import java.util.ArrayList; +import java.util.Iterator; +import java.util.List; +import java.util.Map; +import org.apache.rocketmq.client.producer.DefaultMQProducer; +import org.apache.rocketmq.common.message.Message; + +public class SplitBatchProducer { + + public static void main(String[] args) throws Exception { + + DefaultMQProducer producer = new DefaultMQProducer("BatchProducerGroupName"); + producer.start(); + + //large batch + String topic = "BatchTest"; + List messages = new ArrayList<>(100 * 1000); + for (int i = 0; i < 100 * 1000; i++) { + messages.add(new Message(topic, "Tag", "OrderID" + i, ("Hello world " + i).getBytes())); + } + + //split the large batch into small ones: + ListSplitter splitter = new ListSplitter(messages); + while (splitter.hasNext()) { + List listItem = splitter.next(); + producer.send(listItem); + } + } + +} + + +class ListSplitter implements Iterator> { + private int sizeLimit = 1000 * 1000; + private final List messages; + private int currIndex; + public ListSplitter(List messages) { + this.messages = messages; + } + @Override public boolean hasNext() { + return currIndex < messages.size(); + } + @Override public List next() { + int nextIndex = currIndex; + int totalSize = 0; + for (; nextIndex < messages.size(); nextIndex++) { + Message message = messages.get(nextIndex); + int tmpSize = message.getTopic().length() + message.getBody().length; + Map properties = message.getProperties(); + for (Map.Entry entry : properties.entrySet()) { + tmpSize += entry.getKey().length() + entry.getValue().length(); + } + tmpSize = tmpSize + 20; //for log overhead + if (tmpSize > sizeLimit) { + //it is unexpected that single message exceeds the sizeLimit + //here just let it go, otherwise it will block the splitting process + if (nextIndex - currIndex == 0) { + //if the next sublist has no element, add this one and then break, otherwise just break + nextIndex++; + } + break; + } + if (tmpSize + totalSize > sizeLimit) { + break; + } else { + totalSize += tmpSize; + } + + } + List subList = messages.subList(currIndex, nextIndex); + currIndex = nextIndex; + return subList; + } + + @Override public void remove() { + throw new UnsupportedOperationException("Not allowed to remove"); + } +}