首页 文章 精选 留言 我的

精选列表

搜索[Agent 集群],共10010篇文章
优秀的个人博客,低调大师

Docker部属Nsq集群

用一了段时间NSQ还是很稳定的。除了稳定,还有一个特别值的说的就是部署非常简单。总想写点什么推荐给大家使用nsq来做一些东西。但是就是因为他太简单易用,文档也比较简单易懂。一直不知道要写啥!!!!! nsq官网:http://nsq.io/ 为了容灾需要对nsqd多机器部属,有了Docker后,快速扩还是很方便的。 部署完后我会用go和c#写一些代码方便大家学习。 准备工作: 》两台服务器:192.168.0.49; 192.168.0.105. 》需要在两台机器上安装好Docker 》两台机器上镜像的拉取 docker pull nsqio/nsq 我们在105上启动lookup, nsqd和客户端都需要连接这个lookup。 docker run --name lookupd -p 4160:4160 -p 4161:4161 nsqio/nsq /nsqlookupd 在105和49上启动nsqd, lookup的地址要写105 docker run --name nsqd -p 4150:4150 -p 4151:4151 nsqio/nsq /nsqd --broadcast-address=192.168.0.105 --lookupd-tcp-address=192.168.0.105:4160 docker run --name nsqd -p 4150:4150 -p 4151:4151 nsqio/nsq /nsqd --broadcast-address=192.168.0.49 --lookupd-tcp-address=192.168.0.105:4160 到了这一步就可以写代码发送和接收信息了。但是还有一个管理系统需要启动一下。nsqadmin docker run --name nsqadmin -p 4171:4171 nsqio/nsq /nsqadmin --lookupd-http-address=192.168.0.105:4161 用浏览器看一下管理端:http://192.168.0.105:4171/nodes。找开Nodes标签里面有两个节点。192.168.0.105 和 192.168.0.49。其他的你可以点开看看。 我用go语言 简单写一个发送信息的例子: go使用的库是go-nsq地址 :github.com/nsqio/go-nsq func main() { config := nsq.NewConfig() // 随便给哪个ip发都可以 //w1, _ := nsq.NewProducer("192.168.0.105:4150", config) w1, _ := nsq.NewProducer("192.168.0.49:4150", config) err1 := w1.Ping() if err1 != nil { log.Fatal("should not be able to ping after Stop()") return } defer w1.Stop() topicName := "publishtest" msgCount := 2 for i := 1; i < msgCount; i++ { err1 := w1.Publish(topicName, []byte("测试测试publis test case")) if err1 != nil { log.Fatal("error") } } } 可以尝试给49和105都发送一次试试。再看一下我们的管理页面: publishtest被ip105和49都发送过。但是还没有channel: 客户端golang代码 package main import ( "fmt" "github.com/nsqio/go-nsq" "log" "os" "os/signal" "strconv" "time" "sync" ) func main() { topicName := "publishtest" msgCount := 2 for i := 0; i < msgCount; i++ { //time.Sleep(time.Millisecond * 20) go readMessage(topicName, i) } //cleanup := make(chan os.Signal, 1) cleanup := make(chan os.Signal) signal.Notify(cleanup, os.Interrupt) fmt.Println("server is running....") quit := make(chan bool) go func() { select { case <- cleanup: fmt.Println("Received an interrupt , stoping service ...") for _, ele := range consumers { ele.StopChan <- 1 ele.Stop() } quit <- true } }() <-quit fmt.Println("Shutdown server....") } type ConsumerHandle struct { q *nsq.Consumer msgGood int } var consumers []*nsq.Consumer = make([]*nsq.Consumer, 0) var mux *sync.Mutex = &sync.Mutex{} func (h *ConsumerHandle) HandleMessage(message *nsq.Message) error { msg := string(message.Body) + " " + strconv.Itoa(h.msgGood) fmt.Println(msg) return nil } func readMessage(topicName string, msgCount int) { defer func() { if err := recover(); err != nil { fmt.Println("error: ", err) } }() config := nsq.NewConfig() config.MaxInFlight = 1000 config.MaxBackoffDuration = 500 * time.Second //q, _ := nsq.NewConsumer(topicName, "ch" + strconv.Itoa(msgCount), config) //q, _ := nsq.NewConsumer(topicName, "ch" + strconv.Itoa(msgCount) + "#ephemeral", config) q, _ := nsq.NewConsumer(topicName, "ch"+strconv.Itoa(msgCount), config) h := &ConsumerHandle{q: q, msgGood: msgCount} q.AddHandler(h) err := q.ConnectToNSQLookupd("192.168.0.105:4161") //err := q.ConnectToNSQDs([]string{"192.168.0.105:4161"}) //err := q.ConnectToNSQD("192.168.0.49:4150") //err := q.ConnectToNSQD("192.168.0.105:4415") if err != nil { fmt.Println("conect nsqd error") log.Println(err) } mux.Lock() consumers = append(consumers, q) mux.Unlock() <-q.StopChan fmt.Println("end....") } 本文转自lpxxn博客园博客,原文链接:http://www.cnblogs.com/li-peng/p/7729174.html,如需转载请自行联系原作者 运行一下,会启动两个终端: 用我们的发送代码发送信息,再看我们的客户端 c#使用的库为NsqSharp.Core地址为: https://github.com/tonyredondo/NsqSharp 简单客户端代码为: class Program { static void Main() { // Create a new Consumer for each topic/channel var consumerCount = 2; var listC = new List<Consumer>(); for (var i = 0; i < consumerCount; i++) { var consumer = new Consumer("publishtest", $"channel{i}" ); consumer.ChangeMaxInFlight(2500); consumer.AddHandler(new MessageHandler()); consumer.ConnectToNsqLookupd("192.168.0.105:4161"); listC.Add(consumer); } var exitEvent = new ManualResetEvent(false); Console.CancelKeyPress += (sender, eventArgs) => { eventArgs.Cancel = true; listC.ForEach(x => x.Stop()); exitEvent.Set(); }; exitEvent.WaitOne(); } } public class MessageHandler : IHandler { /// <summary>Handles a message.</summary> public void HandleMessage(IMessage message) { string msg = Encoding.UTF8.GetString(message.Body); Console.WriteLine(msg); } /// <summary> /// Called when a message has exceeded the specified <see cref="Config.MaxAttempts"/>. /// </summary> /// <param name="message">The failed message.</param> public void LogFailedMessage(IMessage message) { // Log failed messages } }

优秀的个人博客,低调大师

Kafka Docker集群搭建

1. Zookeeper下载 http://apache.org/dist/zookeeper/ http://mirrors.hust.edu.cn/apache/zookeeper/zookeeper-3.4.9/zookeeper-3.4.9.tar.gz 2.Kafka下载 http://apache.org/dist/kafka/ http://mirror.bit.edu.cn/apache/kafka/0.10.0.0/kafka_2.10-0.10.0.0.tgz 3,Zookeeper的Docker镜像制作 可以直接使用Zookeeper镜像 docker pull zookeeper 也可以自己使用Dockerfile制作 #基础镜像使用java,这样可以免于设置java环境 FROM java #作者 MAINTAINER Bonker <bonker@foxmail.com> #定义工作目录 ENV WORK_PATH /opt/zkcluster/zkconf #定义zookeeper文件夹名称 ENV ZOOKEEPER_PACKAGE_NAME zookeeper-3.4.9 #创建工作目录 RUN mkdir -p $WORK_PATH #把zookeeper压缩文件复制到工作目录 COPY ./$ZOOKEEPER_PACKAGE_NAME.tar.gz $WORK_PATH/ #解压缩 RUN tar -xvf $WORK_PATH/$ZOOKEEPER_PACKAGE_NAME.tar.gz -C $WORK_PATH/ #删除压缩文件 RUN rm $WORK_PATH/$ZOOKEEPER_PACKAGE_NAME.tar.gz #给shell赋予执行权限 RUN chmod a+x $WORK_PATH/$ZOOKEEPER_PACKAGE_NAME/bin/zkServer.sh RUN mv $WORK_PATH/$ZOOKEEPER_PACKAGE_NAME/conf/zoo_sample.cfg $WORK_PATH/$ZOOKEEPER_PACKAGE_NAME/conf/zoo.cfg ENTRYPOINT $WORK_PATH/$ZOOKEEPER_PACKAGE_NAME/bin/zkServer.sh start-foreground #ENTRYPOINT ["/bin/sh -c",$WORK_PATH/$ZOOKEEPER_PACKAGE_NAME/bin/zkServer.sh] #CMD [start-foreground] 构建zookeeper docker build -t bonker/zookeeper:3.4.9 . 运行docker docker run -d -p 2181:2181 --name zookeeper3.4.9 bonker/zookeeper:3.4.9 4,Kafka的Docker镜像制作 编写执行脚本kafkaStart.sh sed -i "s/broker.id=0/broker.id=$BROKER_ID/g" $WORK_PATH/$KAFKA_PACKAGE_NAME/config/server.properties $WORK_PATH/$KAFKA_PACKAGE_NAME/bin/kafka-server-start.sh $WORK_PATH/$KAFKA_PACKAGE_NAME/config/server.properties 编写Dockerfile #基础镜像使用java,这样可以免于设置java环境 FROM java #作者 MAINTAINER Bonker <bonker@foxmail.com> #定义工作目录 ENV WORK_PATH /usr/local/work #定义kafka文件夹名称 ENV KAFKA_PACKAGE_NAME kafka_2.10-0.10.0.0 #创建工作目录 RUN mkdir -p $WORK_PATH #把kafka压缩文件复制到工作目录 COPY ./$KAFKA_PACKAGE_NAME.tgz $WORK_PATH/ #把kafka执行脚本复制到工作目录 COPY ./kafkaStart.sh $WORK_PATH/ #给shell赋予执行权限 RUN chmod a+x $WORK_PATH/kafkaStart.sh #解压缩 RUN tar -xvf $WORK_PATH/$KAFKA_PACKAGE_NAME.tgz -C $WORK_PATH/ #删除压缩文件 RUN rm $WORK_PATH/$KAFKA_PACKAGE_NAME.tgz #执行sed命令修改文件,将连接zk的ip改为link参数对应的zookeeper容器的别名 RUN sed -i 's/zookeeper.connect=localhost:2181/zookeeper.connect=zkhost:2181/g' $WORK_PATH/$KAFKA_PACKAGE_NAME/config/server.properties #执行sed命令修改文件,改变brokerId #RUN sed -i "s/broker.id=0/broker.id=$BROKER_ID/g" $WORK_PATH/$KAFKA_PACKAGE_NAME/config/server.properties #CMD $WORK_PATH/$KAFKA_PACKAGE_NAME/bin/kafka-server-start.sh $WORK_PATH/$KAFKA_PACKAGE_NAME/config/server.properties CMD $WORK_PATH/kafkaStart.sh 构建kafka docker build -t bonker/kafka:2.10-0.10.0.0 . 5,安装docker-compose 安装python-pip yum -y install epel-release yum -y install python-pip 安装docker-compose pip install docker-compose 待安装完成后,执行查询版本的命令,即可安装docker-compose docker-compose version 6,运行Dokcer容器 编写docker-compose.yml version: '2' services: zk_server: image: bonker/zookeeper:3.4.9 restart: always kafka_server: image: bonker/kafka:2.10-0.10.0.0 ports: - "9091:9092" environment: BROKER_ID: 1 links: - zk_server:zkhost restart: always message_producer: image: bonker/kafka:2.10-0.10.0.0 ports: - "9092:9092" environment: BROKER_ID: 2 links: - zk_server:zkhost restart: always message_consumer: image: bonker/kafka:2.10-0.10.0.0 ports: - "9093:9092" environment: BROKER_ID: 3 links: - zk_server:zkhost restart: always 启动容器 现在打开终端,在docker-compose.yml所在目录下执行docker-compose up -d,即可启动所有容器 作者:Bonker 出处:http://www.cnblogs.com/Bonker QQ:519841366 "> 本页版权归作者和博客园所有,欢迎转载,但未经作者同意必须保留此段声明, 且在文章页面明显位置给出原文链接,否则保留追究法律责任的权利

资源下载

更多资源
Mario

Mario

马里奥是站在游戏界顶峰的超人气多面角色。马里奥靠吃蘑菇成长,特征是大鼻子、头戴帽子、身穿背带裤,还留着胡子。与他的双胞胎兄弟路易基一起,长年担任任天堂的招牌角色。

腾讯云软件源

腾讯云软件源

为解决软件依赖安装时官方源访问速度慢的问题,腾讯云为一些软件搭建了缓存服务。您可以通过使用腾讯云软件源站来提升依赖包的安装速度。为了方便用户自由搭建服务架构,目前腾讯云软件源站支持公网访问和内网访问。

Rocky Linux

Rocky Linux

Rocky Linux(中文名:洛基)是由Gregory Kurtzer于2020年12月发起的企业级Linux发行版,作为CentOS稳定版停止维护后与RHEL(Red Hat Enterprise Linux)完全兼容的开源替代方案,由社区拥有并管理,支持x86_64、aarch64等架构。其通过重新编译RHEL源代码提供长期稳定性,采用模块化包装和SELinux安全架构,默认包含GNOME桌面环境及XFS文件系统,支持十年生命周期更新。

WebStorm

WebStorm

WebStorm 是jetbrains公司旗下一款JavaScript 开发工具。目前已经被广大中国JS开发者誉为“Web前端开发神器”、“最强大的HTML5编辑器”、“最智能的JavaScript IDE”等。与IntelliJ IDEA同源,继承了IntelliJ IDEA强大的JS部分的功能。

用户登录
用户注册