salt分发后,主动将已完成的任务数据推送到redis中,使用redis的生产者模式,进行消息传送
#coding=utf-8 import fnmatch,json,logging import salt.config import salt.utils.event from salt.utils.redis import RedisPool import sys,os,datetime,random import multiprocessing,threading from joi.utils.gobsAPI import PostWeb logger = logging.getLogger(__name__) opts = salt.config.client_config('/data/salt/saltstack/etc/salt/master') r_conn = RedisPool(opts.get('redis_db')).getConn() lock = threading.Lock() class RedisQueueDaemon(object): ''' redis 队列监听器 ''' def __init__(self,r_conn): self.r_conn = r_conn #redis 连接实例 self.task_queue = 'task:prod:queue' #任务消息队列 def listen_task(self): ''' 监听主函数 ''' while True: queue_item = self.r_conn.blpop(self.task_queue,0)[1] print "queue get",queue_item #self.run_task(queue_item) t = threading.Thread(target=self.run_task,args=(queue_item,)) t.start() def run_task(self,info): ''' 执行操作函数 ''' lock.acquire() info = json.loads(info) if info['type'] == 'pushTaskData': task_data = self.getTaskData(info['jid']) task_data = json.loads(task_data) if task_data else [] logger.info('获取缓存数据:%s' % task_data) if task_data: if self.sendTaskData2bs(task_data): task_data = [] self.setTaskData(info['jid'], task_data) elif info['type'] == 'setTaskState': self.setTaskState(info['jid'],info['state'],info['message']) elif info['type'] == 'setTaskData': self.setTaskData(info['jid'], info['data']) lock.release() def getTaskData(self,jid): return self.r_conn.hget('task:'+jid,'data') def setTaskData(self,jid,data): self.r_conn.hset('task:'+jid,'data',json.dumps(data)) def sendTaskData2bs(self,task_data): logger.info('发送任务数据到后端...') logger.info(task_data) if task_data: p = PostWeb('/jgapi/verify',task_data,'pushFlowTaskData') result = p.postRes() print result if result['code']: logger.info('发送成功!') return True else: logger.error('发送失败!') return False else: return True def setTaskState(self,jid,state,message=''): logger.info('到后端设置任务【%s】状态' % str(jid)) p = PostWeb('/jgapi/verify',{'code':jid,'state':'success','message':message},'setTaskState') result = p.postRes() if result['code']: logger.info('设置任务【%s】状态成功!' % str(jid)) return True,result else: logger.error('设置任务【%s】状态失败!' % str(jid)) return result def salt_job_listener(): ''' salt job 监听器 ''' sevent = salt.utils.event.get_event( 'master', sock_dir=opts['sock_dir'], transport=opts['transport'], opts=opts) while True: ret = sevent.get_event(full=True) if ret is None: continue if fnmatch.fnmatch(ret['tag'], 'salt/job/*/ret/*'): task_key = 'task:'+ret['data']['jid'] task_state = r_conn.hget(task_key,'state') task_data = r_conn.hget(task_key,'data') if task_state: jid_data = { 'code':ret['data']['jid'], 'project_id':settings.SALT_MASTER_OPTS['project_id'], 'serverip':ret['data']['id'], 'returns':ret['data']['return'], 'name':ret['data']['id'], 'state':'success' if ret['data']['success'] else 'failed', } task_data = json.loads(task_data) if task_data else [] task_data.append(jid_data) logger.info("新增数据:%s" % json.dumps(task_data)) r_conn.lpush('task:prod:queue',json.dumps({'type':'setTaskData','jid':ret['data']['jid'],'data':task_data})) #r_conn.hset(task_key,'data',json.dumps(task_data)) if task_state == 'running': if len(task_data)>=1: logger.info('新增消息到队列:pushTaskData') r_conn.lpush('task:prod:queue',json.dumps({'jid':ret['data']['jid'],'type':'pushTaskData'})) else: logger.info('任务{0}完成,发送剩下的数据到后端...'.format(task_key)) logger.info('新增消息到队列:pushTaskData') r_conn.lpush('task:prod:queue',json.dumps({'jid':ret['data']['jid'],'type':'pushTaskData'})) print datetime.datetime.now() def run(): print 'start redis product queue listerner...' logger.info('start redis product queue listerner...') multiprocessing.Process(target=RedisQueueDaemon(r_conn).listen_task,args=()).start() print 'start salt job listerner...' logger.info('start salt job listerner...') multiprocessing.Process(target=salt_job_listener,args=()).start() ''' p=multiprocessing.Pool(2) print 'start redis product queue listerner...' p.apply_async(redis_queue_listenr,()) print 'start salt job listerner...' p.apply_async(salt_job_listener,()) p.close() p.join() '''
以上这篇python 监听salt job状态,并任务数据推送到redis中的方法就是小编分享给大家的全部内容了,希望能给大家一个参考,也希望大家多多支持。
免责声明:本站文章均来自网站采集或用户投稿,网站不提供任何软件下载或自行开发的软件!
如有用户或公司发现本站内容信息存在侵权行为,请邮件告知! 858582#qq.com
暂无“python 监听salt job状态,并任务数据推送到redis中的方法”评论...
稳了!魔兽国服回归的3条重磅消息!官宣时间再确认!
昨天有一位朋友在大神群里分享,自己亚服账号被封号之后居然弹出了国服的封号信息对话框。
这里面让他访问的是一个国服的战网网址,com.cn和后面的zh都非常明白地表明这就是国服战网。
而他在复制这个网址并且进行登录之后,确实是网易的网址,也就是我们熟悉的停服之后国服发布的暴雪游戏产品运营到期开放退款的说明。这是一件比较奇怪的事情,因为以前都没有出现这样的情况,现在突然提示跳转到国服战网的网址,是不是说明了简体中文客户端已经开始进行更新了呢?
更新动态
2024年11月26日
2024年11月26日
- 凤飞飞《我们的主题曲》飞跃制作[正版原抓WAV+CUE]
- 刘嘉亮《亮情歌2》[WAV+CUE][1G]
- 红馆40·谭咏麟《歌者恋歌浓情30年演唱会》3CD[低速原抓WAV+CUE][1.8G]
- 刘纬武《睡眠宝宝竖琴童谣 吉卜力工作室 白噪音安抚》[320K/MP3][193.25MB]
- 【轻音乐】曼托凡尼乐团《精选辑》2CD.1998[FLAC+CUE整轨]
- 邝美云《心中有爱》1989年香港DMIJP版1MTO东芝首版[WAV+CUE]
- 群星《情叹-发烧女声DSD》天籁女声发烧碟[WAV+CUE]
- 刘纬武《睡眠宝宝竖琴童谣 吉卜力工作室 白噪音安抚》[FLAC/分轨][748.03MB]
- 理想混蛋《Origin Sessions》[320K/MP3][37.47MB]
- 公馆青少年《我其实一点都不酷》[320K/MP3][78.78MB]
- 群星《情叹-发烧男声DSD》最值得珍藏的完美男声[WAV+CUE]
- 群星《国韵飘香·贵妃醉酒HQCD黑胶王》2CD[WAV]
- 卫兰《DAUGHTER》【低速原抓WAV+CUE】
- 公馆青少年《我其实一点都不酷》[FLAC/分轨][398.22MB]
- ZWEI《迟暮的花 (Explicit)》[320K/MP3][57.16MB]