hive-ddl-gen.py 12 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317
  1. #!/usr/bin/env /usr/bin/python3
  2. # -*- coding:utf-8 -*-
  3. """
  4. Hive DDL 生成器(raw / ods 双层)。
  5. **仅支持 PG 源**:reader.dataSource 必须是 `postgresql/{env}-{instance}`
  6. 形式;mysql 等其他源由 dw_base.io.db.postgresql 的 _resolve_datasource
  7. 直接 NotImplementedError。
  8. 输入 sync ini,从 PG 抽字段类型 + 中文注释:
  9. - raw 层:按 reader.column 顺序渲染全字段 STRING + dt STRING 分区 + ORC + EXTERNAL
  10. - ods 层:按 reader.column 顺序应用 conf/pg-to-hive-type.ini 类型映射,
  11. 末尾加 is_deleted BOOLEAN 软删归一字段,dt STRING 分区 + ORC + EXTERNAL;
  12. 不加 etl_time / src_sys / src_tbl 技术字段(详见 ADR-06)
  13. 写到 stdout(传 -o 时额外落盘 {table_name}_create.sql)。
  14. CLI:
  15. python3 bin/hive-ddl-gen.py -l {raw|ods} -ini jobs/raw/{域}/{table}.ini [-o [DIR]]
  16. 参数:
  17. -l 层级(raw / ods,必填)
  18. -ini sync ini 路径(按项目根解析相对路径,与项目其他 bin 入口一致)
  19. -o 输出目录(任意三态 stdout 始终打印;不传仅 stdout;不带值额外落盘
  20. workspace/{yyyymmdd}/;带值额外落盘 <DIR>/)
  21. 表名由 writer.path 末两段反推(path 末段必须是 dt=... 占位);
  22. ods 表名 = raw 表名首段 'raw_' 替换为 'ods_'。
  23. """
  24. import argparse
  25. import importlib.util
  26. import os
  27. import sys
  28. from configparser import ConfigParser
  29. from datetime import datetime
  30. project_root = os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
  31. sys.path.append(project_root)
  32. from dw_base.io.db import postgresql as pgdb
  33. def _load_sync_gen():
  34. """复用 datax-sync-template-gen 的 query_columns_full(脚本名含连字符,importlib 加载)。"""
  35. spec = importlib.util.spec_from_file_location(
  36. 'datax_sync_template_gen',
  37. os.path.join(project_root, 'bin', 'datax-sync-template-gen.py'),
  38. )
  39. mod = importlib.util.module_from_spec(spec)
  40. spec.loader.exec_module(mod)
  41. return mod
  42. SYNC_GEN = _load_sync_gen()
  43. WORKSPACE_DEFAULT = os.path.join(
  44. project_root, 'workspace', datetime.now().strftime('%Y%m%d'),
  45. )
  46. def _resolve_to_project_root(path):
  47. if os.path.isabs(path):
  48. return path
  49. return os.path.join(project_root, path)
  50. def parse_sync_ini(path):
  51. """解析 sync ini,提取 raw DDL 渲染所需 5 项。"""
  52. if not os.path.isfile(path):
  53. raise FileNotFoundError('sync ini 不存在: ' + path)
  54. cp = ConfigParser()
  55. cp.read(path, encoding='utf-8')
  56. if not cp.has_section('reader'):
  57. raise KeyError('sync ini 缺 [reader] 段: ' + path)
  58. if not cp.has_section('writer'):
  59. raise KeyError('sync ini 缺 [writer] 段: ' + path)
  60. ds_ref = cp.get('reader', 'dataSource').strip()
  61. schema_table = cp.get('reader', 'table').strip()
  62. columns_str = cp.get('reader', 'column').strip()
  63. writer_path = cp.get('writer', 'path').strip()
  64. if '.' not in schema_table:
  65. raise ValueError('reader.table 必须 schema.table 格式: ' + schema_table)
  66. schema, table = schema_table.split('.', 1)
  67. columns = [c.strip() for c in columns_str.split(',') if c.strip()]
  68. if not columns:
  69. raise ValueError('reader.column 为空')
  70. return {
  71. 'ds_ref': ds_ref,
  72. 'schema': schema,
  73. 'table': table,
  74. 'columns': columns,
  75. 'writer_path': writer_path,
  76. }
  77. def reverse_table_name(writer_path):
  78. """从 writer.path 反推 Hive 表名。
  79. path 形如 /user/hive/warehouse/raw.db/{table_name}/dt=${dt}/
  80. 末段必须是 dt=... 占位,倒数第二段即表名。
  81. """
  82. p = writer_path.rstrip('/')
  83. parts = p.rsplit('/', 1)
  84. if len(parts) != 2 or not parts[1].startswith('dt='):
  85. raise ValueError(
  86. 'writer.path 末段必须是 dt=...,无法反推表名: ' + writer_path)
  87. return parts[0].rsplit('/', 1)[-1]
  88. def fetch_column_comments(ds_ref, schema, table):
  89. """连 PG 拿 schema.table 的 {字段名: 中文注释} dict。
  90. 复用 sync-template-gen 的 datasource 解析与 pg_catalog 查询,不另起一套。
  91. """
  92. rows = _fetch_pg_column_rows(ds_ref, schema, table)
  93. return {name: (comment or '') for _, name, comment, _, _ in rows}
  94. def fetch_column_full_rows(ds_ref, schema, table):
  95. """连 PG 拿 schema.table 的全字段 [(attnum, attname, comment, pg_type, pk_flag), ...]。
  96. ods 渲染需要 pg_type,raw 只用 comment。本函数返回原始 rows 给 ods 用。
  97. """
  98. return _fetch_pg_column_rows(ds_ref, schema, table)
  99. def _fetch_pg_column_rows(ds_ref, schema, table):
  100. conn = pgdb.connect(ds_ref)
  101. try:
  102. return SYNC_GEN.query_columns_full(conn, schema, table)
  103. finally:
  104. conn.close()
  105. def normalize_pg_type(pg_type):
  106. """PG type → 映射 conf 查询用的 normalized key。
  107. 规则:
  108. - 小写 + 去首尾空格
  109. - 去括号参数:'numeric(12,2)' → 'numeric','character varying(64)' → 'character varying'
  110. - 去时区后缀:'timestamp(6) without time zone' → 'timestamp'
  111. """
  112. t = pg_type.lower().strip()
  113. if '(' in t and ')' in t:
  114. before = t[:t.index('(')].strip()
  115. after = t[t.index(')') + 1:].strip()
  116. t = (before + ' ' + after).strip()
  117. for suffix in ('without time zone', 'with time zone'):
  118. if t.endswith(suffix):
  119. t = t[:-len(suffix)].strip()
  120. return t
  121. def load_type_mapping(conf_path):
  122. """读 conf/pg-to-hive-type.ini 的 [mapping] 段,返回 {normalized_pg_type: hive_type}。"""
  123. if not os.path.isfile(conf_path):
  124. raise FileNotFoundError('类型映射 conf 不存在: ' + conf_path)
  125. cp = ConfigParser()
  126. cp.read(conf_path, encoding='utf-8')
  127. if not cp.has_section('mapping'):
  128. raise KeyError('类型映射 conf 缺 [mapping] 段: ' + conf_path)
  129. return dict(cp.items('mapping'))
  130. def map_pg_to_hive(pg_type, type_mapping):
  131. """PG 字段类型映射到 Hive 类型;未命中报错让人显式补规则。"""
  132. key = normalize_pg_type(pg_type)
  133. if key not in type_mapping:
  134. raise KeyError(
  135. "PG 类型 '{}'(normalized '{}')不在 conf/pg-to-hive-type.ini 映射表,"
  136. "需显式补规则".format(pg_type, key))
  137. return type_mapping[key]
  138. def reverse_ods_table_name(raw_table_name):
  139. """raw_xxx → ods_xxx;首段必须是 'raw_'。"""
  140. if not raw_table_name.startswith('raw_'):
  141. raise ValueError("raw 表名首段必须是 'raw_': " + raw_table_name)
  142. return 'ods_' + raw_table_name[len('raw_'):]
  143. def render_raw_ddl(table_name, columns, comment_dict):
  144. """渲染 raw 层 DDL:全字段 STRING + dt STRING 分区 + ORC + EXTERNAL。
  145. 字段顺序严格按 columns(已是 sync ini reader.column 裁剪后顺序);
  146. 字段注释从 comment_dict 按字段名查,缺失留空字符串。
  147. """
  148. today = datetime.now().strftime('%Y-%m-%d')
  149. width = max(len(c) for c in columns) + 4
  150. lines = [
  151. '-- 作者:<TODO>',
  152. '-- 日期:' + today,
  153. '-- 工单:<TODO>',
  154. '-- 目的:<TODO>',
  155. '-- 状态:[待执行]',
  156. '-- 备注:<TODO>',
  157. '',
  158. 'DROP TABLE IF EXISTS raw.' + table_name + ';',
  159. '',
  160. 'CREATE EXTERNAL TABLE IF NOT EXISTS raw.' + table_name + ' (',
  161. ]
  162. last_idx = len(columns) - 1
  163. for i, col in enumerate(columns):
  164. comma = ',' if i < last_idx else ''
  165. comment = comment_dict.get(col, '').replace("'", "''")
  166. lines.append(" {col:<{w}}STRING COMMENT '{comment}'{comma}".format(
  167. col=col, w=width, comma=comma, comment=comment))
  168. lines.extend([
  169. ')',
  170. "COMMENT '<TODO>'",
  171. 'PARTITIONED BY (dt STRING)',
  172. 'STORED AS ORC',
  173. "LOCATION '/user/hive/warehouse/raw.db/" + table_name + "';",
  174. '',
  175. ])
  176. return '\n'.join(lines)
  177. def render_ods_ddl(raw_table_name, columns, full_rows, type_mapping):
  178. """渲染 ods 层 DDL:typed 字段 + is_deleted 归一 + dt 分区 + ORC + EXTERNAL。
  179. full_rows: [(attnum, attname, comment, pg_type, pk_flag), ...] 来自 query_columns_full
  180. columns: sync ini reader.column 裁剪后字段列表(与 full_rows 字段名子集对齐)
  181. type_mapping: load_type_mapping 返回的 {normalized_pg_type: hive_type}
  182. 字段顺序按 columns(不依赖 full_rows attnum);缺类型 / 缺注释报错。
  183. 末尾加 is_deleted BOOLEAN 软删归一字段(注释固定)。
  184. """
  185. today = datetime.now().strftime('%Y-%m-%d')
  186. ods_table_name = reverse_ods_table_name(raw_table_name)
  187. width = max(len(c) for c in columns + ['is_deleted']) + 4
  188. by_name = {r[1]: (r[3], r[2] or '') for r in full_rows}
  189. missing = [c for c in columns if c not in by_name]
  190. if missing:
  191. raise KeyError('reader.column 中字段 PG 元数据缺失: ' + ','.join(missing))
  192. lines = [
  193. '-- 作者:<TODO>',
  194. '-- 日期:' + today,
  195. '-- 工单:<TODO>',
  196. '-- 目的:<TODO>',
  197. '-- 状态:[待执行]',
  198. '-- 备注:<TODO>',
  199. '',
  200. 'DROP TABLE IF EXISTS ods.' + ods_table_name + ';',
  201. '',
  202. 'CREATE EXTERNAL TABLE IF NOT EXISTS ods.' + ods_table_name + ' (',
  203. ]
  204. type_width = max(len(map_pg_to_hive(by_name[c][0], type_mapping)) for c in columns)
  205. type_width = max(type_width, len('BOOLEAN')) + 2
  206. for col in columns:
  207. pg_type, comment = by_name[col]
  208. hive_type = map_pg_to_hive(pg_type, type_mapping)
  209. comment = comment.replace("'", "''")
  210. lines.append(" {col:<{w}}{ht:<{tw}}COMMENT '{c}',".format(
  211. col=col, w=width, ht=hive_type, tw=type_width, c=comment))
  212. lines.append(" {col:<{w}}{ht:<{tw}}COMMENT '软删除归一(CASE WHEN del_* THEN TRUE)'".format(
  213. col='is_deleted', w=width, ht='BOOLEAN', tw=type_width))
  214. lines.extend([
  215. ')',
  216. "COMMENT '<TODO>'",
  217. 'PARTITIONED BY (dt STRING)',
  218. 'STORED AS ORC',
  219. "LOCATION '/user/hive/warehouse/ods.db/" + ods_table_name + "';",
  220. '',
  221. ])
  222. return '\n'.join(lines)
  223. def main():
  224. parser = argparse.ArgumentParser(
  225. prog='hive-ddl-gen',
  226. description='Hive DDL 生成器(raw / ods 双层)',
  227. )
  228. parser.add_argument('-l', required=True, choices=['raw', 'ods'],
  229. metavar='LAYER',
  230. help='层级(raw / ods,必填)')
  231. parser.add_argument('-ini', required=True, metavar='PATH',
  232. help='sync ini 路径(按项目根解析相对路径)')
  233. parser.add_argument('-o', nargs='?', const=WORKSPACE_DEFAULT, default=None,
  234. metavar='DIR',
  235. help='输出目录(任意三态 stdout 始终打印;不传仅 stdout;不带值额外落盘 workspace/{yyyymmdd}/;带值额外落盘 <DIR>/)')
  236. args = parser.parse_args()
  237. ini_path = _resolve_to_project_root(args.ini)
  238. spec = parse_sync_ini(ini_path)
  239. table_name = reverse_table_name(spec['writer_path'])
  240. if args.l == 'raw':
  241. comment_dict = fetch_column_comments(
  242. spec['ds_ref'], spec['schema'], spec['table'])
  243. ddl = render_raw_ddl(table_name, spec['columns'], comment_dict)
  244. out_table_name = table_name
  245. else:
  246. full_rows = fetch_column_full_rows(
  247. spec['ds_ref'], spec['schema'], spec['table'])
  248. type_mapping = load_type_mapping(
  249. os.path.join(project_root, 'conf', 'pg-to-hive-type.ini'))
  250. ddl = render_ods_ddl(table_name, spec['columns'], full_rows, type_mapping)
  251. out_table_name = reverse_ods_table_name(table_name)
  252. sys.stdout.write(ddl)
  253. if args.o is not None:
  254. os.makedirs(args.o, exist_ok=True)
  255. out_path = os.path.join(args.o, out_table_name + '_create.sql')
  256. with open(out_path, 'w', encoding='utf-8') as f:
  257. f.write(ddl)
  258. print('已写入: ' + out_path, file=sys.stderr)
  259. if __name__ == '__main__':
  260. main()