import datetime import json from MyUtils.ZillionUtil import ZillionUtil from MyUtils.MysqlUtils import MysqlUtils from dateutil.relativedelta import relativedelta INSERT_SQL = "replace into %s.%s(project_id,date,energy_cooling,energy_heating,energy_ac_terminal,energy_light,energy_others,create_time,update_time) values " DELETE_SQL = "DELETE FROM `energy_week_day` WHERE `project_id` = '%s' AND `date` >= '%s' AND `date` < '%s' " def datetime_now(): datetime_now = datetime.datetime.now().strftime("%Y%m%d%H%M%S") return datetime_now # 获取能耗数据 def get_data_time(hbase_database, hbase_table, building, from_time, to_time): Criteria = { "building": building, "energyModelSign": building + "EM01", "energyModelNodeSign": "EI1001", "data_time": { "$gte": from_time, "$lt": to_time } } datas = zillionUtil.select(hbase_database, hbase_table, Criteria) return datas # mysql插入数据 def insert_mysql(sqls,building,my_table): print("%s,开始往mysql插入%s的%s数据..."%(datetime_now(),building,my_table)) for i in range(0, len(sqls), 1000): sqlranges = sqls[i:i + 1000] sqlranges = INSERT_SQL % (my_database,my_table) + ",".join(sqlranges) MysqlUtil.update(sqlranges) print("%s,mysql数据%s,项目%s插入成功,合计%s条..." % (datetime_now(),my_table,building,len(sqls))) with open("config.json", "r") as f: data = json.load(f) hbase_database = data["metadata"]["database"] url = data["metadata"]["url"] building = data["building"] mysql = data["mysql"] my_database = mysql["database"] now_month = datetime.date.today().strftime("%Y%m") last_month = (datetime.date.today() - relativedelta(months=1)).strftime("%Y%m") # #连接hbase zillionUtil = ZillionUtil(url) # #连接hbase MysqlUtil = MysqlUtils(**mysql) datas = get_data_time(hbase_database, "data_energydata_1d", building, last_month, now_month) if datas: project_id = "Pj" + building # 删除上月数据 print("%s,开始删除%s的数据..."%(datetime_now(),last_month)) delete_sql = DELETE_SQL% (project_id,last_month,now_month) MysqlUtil.update(delete_sql) sqls = [] for i in datas: dt = i["data_time"][0:8] data_value = i["data_value"] sql = "('%s','%s','0','0','0','0','%s','%s','%s')" % (project_id, dt, data_value, datetime_now(),datetime_now()) sqls.append(sql) insert_mysql(sqls,building,"energy_week_day") else: print("%s,没有查询到数据...")