metaq总结

ywzqwwt(25)
Published in
#metaq
Words
417
Reading
2 min
Listen
Play
8y

1.使用的包直接导的metaq软件中的lib里面的所有
2.主要写的生产者和消费者
3.另外特别注意需要注册 topic 在metaq/conf下面的server.ini

上代码:生产者
package examples;

/*

  • 生产者 本次用的控制台输入方式。
    */

import java.io.BufferedReader;
import java.io.InputStreamReader;

import com.taobao.metamorphosis.Message;
import com.taobao.metamorphosis.client.MessageSessionFactory;
import com.taobao.metamorphosis.client.MetaClientConfig;
import com.taobao.metamorphosis.client.MetaMessageSessionFactory;
import com.taobao.metamorphosis.client.producer.MessageProducer;
import com.taobao.metamorphosis.client.producer.SendResult;
import com.taobao.metamorphosis.utils.ZkUtils.ZKConfig;

public class Producer {
public static void main(String[] args) throws Exception {
final MetaClientConfig metaClientConfig = new MetaClientConfig();
final ZKConfig zkConfig = new ZKConfig();
//设置zookeeper地址
zkConfig.zkConnect = "197.3.155.125:52181,197.3.155.126:52181,197.3.155.127:52181";
metaClientConfig.setZkConfig(zkConfig);
// New session factory,强烈建议使用单例
MessageSessionFactory sessionFactory = new MetaMessageSessionFactory(metaClientConfig);
// create producer,强烈建议使用单例
MessageProducer producer = sessionFactory.createProducer();
// publish topic 这需要在metaq里面注册test
final String topic = "test";
producer.publish(topic);

    BufferedReader reader = new BufferedReader(new InputStreamReader(System.in));
    String line = null;
    while ((line = reader.readLine()) != null) {
        // send message
        SendResult sendResult = producer.sendMessage(new Message(topic, line.getBytes()));
        // check result
        if (!sendResult.isSuccess()) {
            System.err.println("Send message failed,error message:" + sendResult.getErrorMessage());
        }
        else {
            System.out.println("Send message successfully,sent to " + sendResult.getPartition());
        }
    }
}

}

消费者
package examples;

/*

  • 消费者 消费从生产者那里接收数据,本次使用控制台输入
    */
    import java.util.concurrent.Executor;

import com.taobao.metamorphosis.Message;
import com.taobao.metamorphosis.client.MessageSessionFactory;
import com.taobao.metamorphosis.client.MetaClientConfig;
import com.taobao.metamorphosis.client.MetaMessageSessionFactory;
import com.taobao.metamorphosis.client.consumer.ConsumerConfig;
import com.taobao.metamorphosis.client.consumer.MessageConsumer;
import com.taobao.metamorphosis.client.consumer.MessageListener;
import com.taobao.metamorphosis.utils.ZkUtils.ZKConfig;

public class Consumer {
public static void main(String[] args) throws Exception {
final MetaClientConfig metaClientConfig = new MetaClientConfig();
final ZKConfig zkConfig = new ZKConfig();
//设置zookeeper地址
zkConfig.zkConnect = "197.3.155.125:52181,197.3.155.126:52181,197.3.155.127:52181";
metaClientConfig.setZkConfig(zkConfig);
// New session factory,强烈建议使用单例
MessageSessionFactory sessionFactory = new MetaMessageSessionFactory(metaClientConfig);
// subscribed topic 这需要在metaq里面注册test
final String topic = "test";
// consumer group
final String group = "meta-example";
// create consumer,强烈建议使用单例
MessageConsumer consumer = sessionFactory.createConsumer(new ConsumerConfig(group));
// subscribe topic
consumer.subscribe(topic, 1024 * 1024, new MessageListener() {

        public void recieveMessages(Message message) {
            System.out.println("Receive message " + new String(message.getData()));
        }


        public Executor getExecutor() {
            // Thread pool to process messages,maybe null.
            return null;
        }
    });
    // complete subscribe
    consumer.completeSubscribe();
}

}

metaq总结 | Ecency