Redis publish and subscribe

linkaiyi
Follow

reciever

# encoding: utf-8
__author__ = 'DDK'
__copyright__ = 'Copyright 2009, DDK'
__version__ = '1.0'

from django.core.management.base import BaseCommand
from interface.models import DailySalesReport, MonthSalesReport, QuarterSalesReport
from interface.views import update_daily_sales, update_quarterly_sales, update_monthly_sales, update_orgchart
import threading
from time import sleep
import redis


r = redis.StrictRedis(host='localhost', port=6379, db=4)
p = r.pubsub()
p.subscribe('update_orgchart', 'update_monthly_sales', 'update_quarterly_sales', 'update_daily_sales')


class Command(BaseCommand):
    help = "Update Daily Sales Table"

    def handle(self, *args, **options):
        sleep(300)
        DailySalesReport.objects.all().delete()
        MonthSalesReport.objects.all().delete()
        QuarterSalesReport.objects.all().delete()
        orgchart_status = True
        daily_sales_status = True
        monthly_sales_status = True
        quarterly_sales_status = True
        orgchart = threading.Thread(target=update_orgchart, daemon=True)
        orgchart.start()
        daily_sales = threading.Thread(target=update_daily_sales, daemon=True)
        daily_sales.start()
        sleep(0.1)
        monthly_sales = threading.Thread(target=update_monthly_sales, daemon=True)
        monthly_sales.start()
        sleep(0.1)
        quarterly_sales = threading.Thread(target=update_quarterly_sales, daemon=True)
        quarterly_sales.start()
        while orgchart_status | daily_sales_status | monthly_sales_status | quarterly_sales_status:
            message = p.get_message()
            # print(message)
            # print('orgchart: {} | daily: {} | monthly: {} | quarterly: {}'.format(str(orgchart_status), str(daily_sales_status), str(monthly_sales_status), str(quarterly_sales_status)))
            if message:
                if orgchart_status and message['channel'] == b'update_orgchart' and message['data'] == b'finished':
                    orgchart_status = False
                if daily_sales_status and message['channel'] == b'update_daily_sales' and message['data'] == b'finished':
                    daily_sales_status = False
                if monthly_sales_status and message['channel'] == b'update_monthly_sales' and message['data'] == b'finished':
                    monthly_sales_status = False
                if quarterly_sales_status and message['channel'] == b'update_quarterly_sales' and message['data'] == b'finished':
                    quarterly_sales_status = False
            sleep(1)

        self.stdout.write('Successfull updated!')

sender

import requests
from organization.models import OrgChartProperty, OrgChart
from interface.models import OrgChartOne, OrgChartTwo, OrgChartThree, MonthSalesReport, QuarterSalesReport, DailySalesReport
from operation.models import PostEmployeeRelation, PostPartnerRelation
import json
from time import sleep
import redis
import urllib3
# Create your views here.

MONTHLY = 'https://www.ace-bdc.com.cn:4243/v1/doc/489186d2-d9a9-4bc1-8741-a4000bb3dadd/object/1104e96d-4b5d-412d-8e21-d23d1a57b26a/data/'
QUARTERLY = 'https://www.ace-bdc.com.cn:4243/v1/doc/489186d2-d9a9-4bc1-8741-a4000bb3dadd/object/b0dcdc50-b8ec-4270-badc-e3a0e8bbe63c/data/'
DAILY = 'https://www.ace-bdc.com.cn:4243/v1/doc/489186d2-d9a9-4bc1-8741-a4000bb3dadd/object/jGgUPe/data/'

urllib3.disable_warnings(urllib3.exceptions.InsecureRequestWarning)

r = redis.StrictRedis(host='localhost', port=6379, db=4)
P = r.pubsub()


def update_orgchart():
    r.publish('update_orgchart', 'started')
    try:
        OrgChartOne.objects.filter().delete()
        OrgChartTwo.objects.all().delete()
        OrgChartThree.objects.all().delete()
        for item in OrgChartProperty.objects.all().exclude(post__name__istartswith='BD'):
            OrgChartOne.objects.create(start_time=item.valid_date, end_time=item.valid_end, org_id=item.post.code,
                                       org_name=item.post.name, is_virtual_org=item.is_virtual, org_type=item.post.get_org_type_display(),
                                       parent_org_id=item.parent.post.code if item.parent else None, last_modified_time=item.modified)

        for item in PostEmployeeRelation.objects.all().exclude(post__name__istartswith='BD'):
            OrgChartTwo.objects.create(org_id=item.post.code if item.post else None, employee_id=item.employee.code if item.employee else None,
                                       start_time=item.valid_date,
                                       end_time=item.valid_end, is_co_manage=item.virtual_host,
                                       last_modified_time=item.modified)

        for item in PostPartnerRelation.objects.all().exclude(post__name__istartswith='BD'):
            OrgChartThree.objects.create(start_time=item.valid_date, end_time=item.valid_end, org_id=item.post.code,
                                         hco_id=item.partner.code, product_id=item.product.code, co_manage_ratio=item.management_index,
                                         last_modified_time=item.modified)
        r.publish('update_orgchart', 'finished')

    except Exception as e:
        r.publish('update_orgchart', str(e))


def update_monthly_sales():
    r.publish('update_monthly_sales', 'started')
    try:
        session_get = requests.session()
        response = session_get.get(MONTHLY + '1', verify=False)
        data = json.loads(response.text)
        page_num = data['page_total']
        obj_list = []
        for i in range(1, page_num+1):
            for item in data['data'][0]['qMatrix']:
                new_obj = MonthSalesReport(
                            rep_name=item[0]['qText'],
                            rep_code=item[1]['qText'],
                            territory_code=item[2]['qText'],
                            product_id=item[3]['qText'],
                            product_name=item[4]['qText'],
                            sales_amount=item[5]['qText'],
                            target_amount=item[6]['qText'],
                            sales_volume=item[7]['qText'],
                            target_volume=item[8]['qText'],
                            org_code=item[9]['qText'],
                            org_name=item[10]['qText'],
                            last_year_amount=item[11]['qText'],
                            datetime=item[12]['qText'],
                        )
                obj_list.append(new_obj)
            MonthSalesReport.objects.bulk_create(obj_list)
            obj_list = []
            if i < page_num:
                sleep(0.05)
                response = session_get.get(MONTHLY + str(i+1), verify=False)
                data = json.loads(response.text)
                # print(MONTHLY + str(i+1))
        r.publish('update_monthly_sales', 'finished')

    except Exception as e:
        r.publish('update_monthly_sales', str(e))


def update_quarterly_sales():
    r.publish('update_quarterly_sales', 'started')
    try:
        session_get = requests.session()
        response = session_get.get(QUARTERLY + '1', verify=False)
        data = json.loads(response.text)
        page_num = data['page_total']
        obj_list = []
        for i in range(1, page_num+1):
            for item in data['data'][0]['qMatrix']:
                new_obj = QuarterSalesReport(
                            rep_name=item[0]['qText'],
                            rep_code=item[1]['qText'],
                            territory_code=item[2]['qText'],
                            product_id=item[3]['qText'],
                            product_name=item[4]['qText'],
                            sales_amount=item[5]['qText'],
                            target_amount=item[6]['qText'],
                            sales_volume=item[7]['qText'],
                            target_volume=item[8]['qText'],
                            org_code=item[9]['qText'],
                            org_name=item[10]['qText'],
                            last_year_amount=item[11]['qText'],
                            datetime=item[12]['qText'],
                        )
                obj_list.append(new_obj)
            QuarterSalesReport.objects.bulk_create(obj_list)
            obj_list = []
            if i < page_num:
                sleep(0.05)
                response = session_get.get(QUARTERLY + str(i+1), verify=False)
                data = json.loads(response.text)
                # print(QUARTERLY + str(i+1))
        r.publish('update_quarterly_sales', 'finished')

    except Exception as e:
        r.publish('update_quarterly_sales', str(e))


def update_daily_sales():
    r.publish('update_daily_sales', 'started')
    try:
        session_get = requests.session()
        response = session_get.get(DAILY + '1', verify=False)
        data = json.loads(response.text)
        page_num = data['page_total']
        obj_list = []
        for i in range(1, page_num+1):
            for item in data['data'][0]['qMatrix']:
                new_obj = DailySalesReport(
                            rep_name=item[0]['qText'],
                            rep_code=item[1]['qText'],
                            territory_code=item[2]['qText'],
                            sales_amount=item[3]['qText'],
                            sales_volume=item[4]['qText'],
                            org_code=item[5]['qText'],
                            org_name=item[6]['qText'],
                            product_id=item[7]['qText'],
                            product_name=item[8]['qText'],
                            datetime=item[9]['qText'],
                        )
                obj_list.append(new_obj)
            DailySalesReport.objects.bulk_create(obj_list)
            obj_list = []
            if i < page_num:
                sleep(0.05)
                response = session_get.get(DAILY + str(i+1), verify=False)
                data = json.loads(response.text)
                # print(DAILY + str(i+1))
        r.publish('update_daily_sales', 'finished')
    except Exception as e:
        r.publish('update_daily_sales', str(e))
Object has 0 attachments

Was this article helpful?

This article is viewed 35 times!

Recent Viewed Articles

0 Comments

Leave a Comment

Support