Browse Source

transfer_data

李莎 2 years ago
parent
commit
c2435ba892
2 changed files with 56 additions and 48 deletions
  1. 10 7
      config.json
  2. 46 41
      start_mysql.py

+ 10 - 7
config.json

@@ -9,16 +9,19 @@
     "port": 9092
   },
   "building": {
-    "id": "1101080259"
+    "id": [
+      "1101110001",
+      "1101080259"
+    ]
   },
   "mysql": {
-    "host": "192.168.100.197",
-    "port": 3306,
+    "host": "bj-cdb-5om355fk.sql.tencentcdb.com",
+    "port": 59790,
     "user": "root",
-    "password": "j5ry0#jZ7vaUt5f4",
-    "database": "saga_dev",
-    "fjd_table": "energy_15_min_fjd_py",
-    "energy_table": "energy_15_min_py",
+    "password": "dwh@2022",
+    "database": "dwh_saga",
+    "fjd_table": "energy_15_min_fjd",
+    "energy_table": "energy_15_min",
     "co2_table": "co2_15_min",
     "pm25_table": "pm25_15_min",
     "temperature_table": "temperature_15_min",

+ 46 - 41
start_mysql.py

@@ -37,11 +37,11 @@ def get_data_time(hbase_database, hbase_table, buildingid, meter, funcid,from_ti
     return datas
 
 #取hbase数据,处理成sql语句
-def hbase_energy_data(points,table):
+def hbase_energy_data(points,table,building):
     sqls = []
     for i in points:
         meter,funcid = i[0],i[1]
-        print("%s开始查询%s至%s的数据 %s-%s"%(table,args.start_time,args.end_time,meter,funcid))
+        print("%s开始查询项目%s的%s至%s的数据 %s-%s"%(table,building,args.start_time,args.end_time,meter,funcid))
         datas = get_data_time(hbase_database,table,building,meter,funcid,args.start_time,args.end_time)
         for i in datas:
             data_time = i["data_time"]
@@ -53,14 +53,14 @@ def hbase_energy_data(points,table):
 
 
 # mysql插入数据
-def insert_mysql(sqls,my_table):
-    print("开始往mysql插入%s数据..."%my_table)
+def insert_mysql(sqls,building,my_table):
+    print("开始往mysql插入%s的%s数据..."%(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)
         mysql_cur.execute(sqlranges)
         conn.commit()
-    print("mysql数据%s插入成功,合计%s条..." % (my_table,len(sqls)))
+    print("mysql数据%s,项目%s插入成功,合计%s条..." % (my_table,building,len(sqls)))
 
 
 end_time = datetime.datetime.now().strftime("%Y%m%d") + "000000"
@@ -81,7 +81,8 @@ with open("config.json", "r") as f:
     data = json.load(f)
     hbase_database = data["metadata"]["database"]
     url = data["metadata"]["url"]
-    building = data["building"]["id"]
+    building_list = data["building"]["id"]
+    print("项目列表:%s"%building_list)
     my_database = data["mysql"]["database"]
     my_fjd_table = data["mysql"]["fjd_table"]
     my_energy_table = data["mysql"]["energy_table"]
@@ -103,44 +104,48 @@ with open("config.json", "r") as f:
 tables = ["fjd_0_near_15min","data_servicedata_15min"]
 # #连接hbase
 zillionUtil = ZillionUtil(url)
-#处理点位表
-pointlist = get_pointlist(hbase_database,"dy_pointlist",building)
-points = []
-for point in pointlist:
-    i = point[1]
-    if i == args.funcid:
-        points.append(point)
 #连接mysql
 conn = pymysql.connect(**mysql)
 mysql_cur = conn.cursor()
-
-#电
-if args.funcid == 10101:
-    sqls = hbase_energy_data(points, "fjd_0_near_15min")
-    insert_mysql(sqls, my_fjd_table)
-    sqls = hbase_energy_data(points, "data_servicedata_15min")
-    insert_mysql(sqls, my_energy_table)
-#CO2
-if args.funcid == 11301:
-    sqls = hbase_energy_data(points, "fjd_0_near_15min")
-    insert_mysql(sqls, my_co2_table)
-#pm2.5
-if args.funcid == 11401:
-    sqls = hbase_energy_data(points, "fjd_0_near_15min")
-    insert_mysql(sqls, my_pm25_table)
-#甲醛
-if args.funcid == 11305:
-    sqls = hbase_energy_data(points, "fjd_0_near_15min")
-    insert_mysql(sqls, my_temperature_table)
-#温度
-if args.funcid == 11101:
-    sqls = hbase_energy_data(points, "fjd_0_near_15min")
-    insert_mysql(sqls, my_hcho_table)
-#湿度
-if args.funcid == 11201:
-    sqls = hbase_energy_data(points, "fjd_0_near_15min")
-    insert_mysql(sqls, my_humidity_table)
-
+for building in building_list:
+    #处理点位表
+    pointlist = get_pointlist(hbase_database,"dy_pointlist",building)
+    if pointlist:
+        points = []
+        for point in pointlist:
+            i = point[1]
+            if i == args.funcid:
+                points.append(point)
+
+
+        #电
+        if args.funcid == 10101:
+            sqls = hbase_energy_data(points, "fjd_0_near_15min",building)
+            insert_mysql(sqls,building, my_fjd_table)
+            sqls = hbase_energy_data(points, "data_servicedata_15min",building)
+            insert_mysql(sqls,building, my_energy_table)
+        #CO2
+        if args.funcid == 11301:
+            sqls = hbase_energy_data(points, "fjd_0_near_15min",building)
+            insert_mysql(sqls,building, my_co2_table)
+        #pm2.5
+        if args.funcid == 11401:
+            sqls = hbase_energy_data(points, "fjd_0_near_15min",building)
+            insert_mysql(sqls,building, my_pm25_table)
+        #甲醛
+        if args.funcid == 11305:
+            sqls = hbase_energy_data(points, "fjd_0_near_15min",building)
+            insert_mysql(sqls,building, my_hcho_table)
+        #温度
+        if args.funcid == 11101:
+            sqls = hbase_energy_data(points, "fjd_0_near_15min",building)
+            insert_mysql(sqls,building, my_temperature_table)
+        #湿度
+        if args.funcid == 11201:
+            sqls = hbase_energy_data(points, "fjd_0_near_15min",building)
+            insert_mysql(sqls, building,my_humidity_table)
+    else:
+        print("未查询到项目%s的点位表..."%building)
 mysql_cur.close()
 conn.close()