start.py 3.8 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798
  1. #!/usr/bin/python3
  2. # -*- coding: utf-8 -*-
  3. import json
  4. from MyUtils.DateUtils import *
  5. from MyUtils.ConfigUtils import ConfigUtils
  6. from MyUtils.ZillionUtil import ZillionUtil
  7. import os,time,datetime
  8. #插入hbase
  9. def put_hbase(hbasedatabase,hbasetable,datas):
  10. for i in range(0, len(datas), 1000):
  11. dataranges = datas[i:i + 1000]
  12. zillionUtil.put(hbasedatabase, hbasetable, dataranges)
  13. #删除hbase数据
  14. def remove_hbase(hbasedatabase,hbasetable,key):
  15. zillionUtil.remove(hbasedatabase, hbasetable, key)
  16. #删除hbase数据
  17. def delete_hbase(hbasedatabase,hbasetable):
  18. Criteria = {"project_id":"3301100002"}
  19. zillionUtil.delete(hbasedatabase, hbasetable, Criteria)
  20. #获取表主键
  21. def get_hbase_key(hbasedatabase, hbasetable):
  22. keys = zillionUtil.table_key(hbasedatabase, hbasetable)
  23. return keys
  24. datetimenow = datetime.datetime.now().strftime("%Y-%m-%d %H:%M:%S")
  25. datetimenow_strp = datetime.datetime.strptime(datetimenow,"%Y-%m-%d %H:%M:%S")
  26. print(type(datetimenow_strp),datetimenow_strp)
  27. config = ConfigUtils("config.xml")
  28. url, hbasedatabase = config.readTop("metadata", ["url", "database"])
  29. sourcepath, targetpath = config.readTop("cmd", ["sourcepath", "targetpath"])
  30. #连接hbase
  31. zillionUtil = ZillionUtil(url)
  32. # #切换到输出目录,清空文件
  33. # os.chdir(targetpath)
  34. # os.system("rm -rf *")
  35. #切换到java程序工作目录,执行导出数据程序
  36. os.chdir(sourcepath)
  37. print("进入%s目录,开始执行java程序"%os.getcwd())
  38. cmd = "java -jar -Dfile.encoding=UTF-8 data-migration.jar"
  39. status = os.system(cmd)
  40. #如果程序执行成功
  41. if status == 0:
  42. print("导出程序执行成功")
  43. #切换到输出目录,判断文件是否是最新文件
  44. os.chdir(targetpath)
  45. print("进入%s目录,检查文件并写入hbase"%os.getcwd())
  46. list = os.listdir(os.getcwd())
  47. for file in list:
  48. #linux 导出程序bug,需要处理下文件名称
  49. # print(i)
  50. # file = i.lstrip("\\")
  51. # cmd_mv = "mv \%s %s"%(i,file)
  52. # print(cmd_mv)
  53. # os.system(cmd_mv)
  54. updatetime = os.path.getmtime(file) #查询文件修改时间
  55. timeArray = time.localtime(updatetime)
  56. otherStyleTime = time.strftime("%Y-%m-%d %H:%M:%S", timeArray)
  57. otherStyleTime_strp = datetime.datetime.strptime(otherStyleTime,"%Y-%m-%d %H:%M:%S")
  58. detal_time = (otherStyleTime_strp - datetimenow_strp).total_seconds()
  59. print("%s文件最新修改时间:%s"%(file,otherStyleTime))
  60. hbasetable = file.split(".")[0]
  61. if int(detal_time) < 10000:
  62. with open(file,"r",encoding = 'utf-8') as fp:
  63. print(fp.name)
  64. if file == "rel_btw_objs.json":
  65. datas_rel = []
  66. for line in fp.readlines():
  67. data = json.loads(line)
  68. datas_rel.append(data)
  69. print("删除%s"%fp.name)
  70. delete_hbase(hbasedatabase,hbasetable)
  71. put_hbase(hbasedatabase, hbasetable, datas_rel)
  72. else:
  73. keys = get_hbase_key(hbasedatabase, hbasetable)
  74. datas = []
  75. for line in fp.readlines():
  76. data = json.loads(line)
  77. hbase_key = {}
  78. for key in keys:
  79. if key in data:
  80. hbase_key[key] = data[key]
  81. datas.append(data)
  82. remove_hbase(hbasedatabase,hbasetable,hbase_key)
  83. time.sleep(0.1)
  84. put_hbase(hbasedatabase,hbasetable,datas)
  85. print("%s写入完成"%hbasetable)
  86. else:
  87. print("检查%s导出文件数据,可能不是最新数据"%file)
  88. else:
  89. print("执行java程序失败")