__init__.py 3.8 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879
  1. #!/usr/bin/env /usr/bin/python3
  2. # -*- coding:utf-8 -*-
  3. import os
  4. import socket
  5. import sys
  6. import time
  7. from dw_base.utils.env_loader import bootstrap_env
  8. # bootstrap_env() 在 driver/边缘节点加载 conf/env.sh(返回 True);Spark executor 上 dw_base 以 zip 分发、
  9. # 不含 conf/,返回 False。下方 driver 专属初始化(读 env 变量 / 路径推断 / banner)整块跳过——
  10. # executor 只需能 import 子模块(如 dw_base.tracking.mask,纯 stdlib)供 Python UDF 运行。
  11. if bootstrap_env():
  12. # HADOOP_CONF_DIR / SPARK_CONF_DIR 由 conf/env.sh 提供默认值经 bootstrap_env 注入(shell 侧 export 优先):
  13. # - HADOOP_CONF_DIR:spark-submit 启动 YARN 校验需要;DataX JVM 不读 classpath conf,HA 由 ini [hadoop_config] 节显式注入
  14. # - SPARK_CONF_DIR:pip pyspark 默认指向自身空 conf/,显式指到集群配置才能加载 hive-site.xml,否则 enableHiveSupport 回落 in-memory metastore
  15. # PYSPARK_*:复用 PYTHON3_PATH,避免与 conf/env.sh 双份硬编码
  16. os.environ.setdefault('PYSPARK_DRIVER_PYTHON', os.environ['PYTHON3_PATH'])
  17. os.environ.setdefault('PYSPARK_PYTHON', os.environ['PYTHON3_PATH'])
  18. os.environ['PYTHONUNBUFFERED'] = 'x'
  19. PROJECT_ROOT_PATH = os.path.abspath(os.path.dirname(os.path.dirname(__file__)))
  20. PROJECT_NAME = os.path.basename(PROJECT_ROOT_PATH)
  21. sys.path.append(PROJECT_ROOT_PATH)
  22. # 公用的Spark UDF文件
  23. COMMON_SPARK_UDF_FILE = 'dw_base/udf/common/spark_common_udf.py'
  24. BANNED_USER = 'root'
  25. RELEASE_USER = os.environ['RELEASE_USER']
  26. USER = os.environ['USER']
  27. HOME = os.environ['HOME']
  28. if USER == BANNED_USER and HOME.startswith('/home'):
  29. USER = os.path.basename(HOME)
  30. HOST = socket.gethostname()
  31. RELEASE_ROOT_DIR = os.environ['RELEASE_ROOT_DIR']
  32. if not PROJECT_ROOT_PATH.startswith(RELEASE_ROOT_DIR) or USER != RELEASE_USER:
  33. DO_RESET: str = '\033[0m'
  34. NORM_RED: str = '\033[0;31m'
  35. NORM_GRN: str = '\033[0;32m'
  36. NORM_YEL: str = '\033[0;33m'
  37. NORM_MGT: str = '\033[0;35m'
  38. NORM_CYN: str = '\033[0;36m'
  39. else:
  40. DO_RESET: str = ''
  41. NORM_RED: str = ''
  42. NORM_GRN: str = ''
  43. NORM_YEL: str = ''
  44. NORM_MGT: str = ''
  45. NORM_CYN: str = ''
  46. IS_RUN_BY_RELEASE_USER = False
  47. LOG_ROOT_DIR = os.environ['LOG_ROOT_DIR']
  48. if USER == RELEASE_USER:
  49. IS_RUN_BY_RELEASE_USER = True
  50. elif USER == BANNED_USER:
  51. ERROR_CODE = 18
  52. print(f'{NORM_MGT}{time.strftime("%Y-%m-%d %H:%M:%S", time.localtime())} '
  53. f'{NORM_RED}Project {NORM_GRN}{PROJECT_NAME} '
  54. f'{NORM_RED}is running by banned user {NORM_GRN}{BANNED_USER}'
  55. f'{NORM_RED}, exit with error code {NORM_GRN}{ERROR_CODE}'
  56. f'{DO_RESET}')
  57. exit(ERROR_CODE)
  58. else:
  59. print(f'{NORM_CYN}{time.strftime("%Y-%m-%d %H:%M:%S", time.localtime())} '
  60. f'{NORM_MGT}Project {NORM_GRN}{PROJECT_NAME} '
  61. f'{NORM_MGT}is running in normal user {NORM_GRN}{USER}')
  62. if PROJECT_ROOT_PATH.startswith(f'{RELEASE_ROOT_DIR}/{PROJECT_NAME}'):
  63. IS_RUN_IN_RELEASE_DIR = True
  64. print(f'{NORM_CYN}{time.strftime("%Y-%m-%d %H:%M:%S", time.localtime())} '
  65. f'{NORM_MGT}Project {NORM_GRN}{PROJECT_NAME} '
  66. f'{NORM_MGT}is running in release dir {NORM_GRN}{RELEASE_ROOT_DIR}/{PROJECT_NAME}')
  67. else:
  68. IS_RUN_IN_RELEASE_DIR = False
  69. print(f'{NORM_CYN}{time.strftime("%Y-%m-%d %H:%M:%S", time.localtime())} '
  70. f'{NORM_MGT}Project {NORM_GRN}{PROJECT_NAME} '
  71. f'{NORM_MGT}is running in normal user dir {NORM_GRN}{PROJECT_ROOT_PATH}')
  72. if not IS_RUN_IN_RELEASE_DIR or USER != RELEASE_USER:
  73. os.system(f'echo -en "{NORM_GRN}"')
  74. os.system(f'echo -en "{DO_RESET}"')