job_config_generator.py 3.7 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192
  1. # -*- coding:utf-8 -*-
  2. import json
  3. import os
  4. from configparser import ConfigParser
  5. from typing import Dict, Optional
  6. from dw_base.datax.datax_constants import *
  7. from dw_base.datax.plugins.plugin_factory import PluginFactory
  8. from dw_base.datax.tuning import load_tuning_conf, merge_speed
  9. from dw_base.utils.file_utils import delete_file, get_abs_path
  10. class JobConfigGenerator(object):
  11. """
  12. 生成 DataX 作业配置文件(json)。
  13. speed 三参走三级合并:L1 conf/datax-tuning.conf < L2 ini [speed] 段 < L3 cli_speed_overrides。
  14. 合并后:channel 写 job.setting.speed(DataX 按 channel 数并发);byte/record 只写
  15. core.transport.channel.speed 当每通道限速——不可再塞进 job.setting.speed,否则并发被压成 1。
  16. """
  17. def __init__(self, base_dir: str, generator_config: str, start_date: str, stop_date: str, output: str,
  18. cli_speed_overrides: Optional[Dict[str, int]] = None):
  19. """
  20. Args:
  21. base_dir: 项目目录
  22. generator_config: DataX 作业配置生成器配置文件路径(.ini)
  23. start_date / stop_date: 内部日期
  24. output: 生成的 DataX json 输出路径
  25. cli_speed_overrides: L3 CLI 覆盖,形如 {'channel': 20, 'byte': None, 'record': None}
  26. """
  27. self.generator_config = get_abs_path(generator_config)
  28. self.base_dir = base_dir
  29. self.start_date = start_date
  30. self.stop_date = stop_date
  31. self.output = output
  32. self.cli_speed_overrides = cli_speed_overrides or {}
  33. self.config_parser = ConfigParser()
  34. self.config_parser.read(self.generator_config)
  35. def get_reader(self):
  36. reader = PluginFactory.get_plugin('reader', self.base_dir, self.config_parser, self.start_date, self.stop_date)
  37. return reader.configure()
  38. def get_writer(self):
  39. writer = PluginFactory.get_plugin('writer', self.base_dir, self.config_parser, self.start_date, self.stop_date)
  40. return writer.configure()
  41. def assemble(self):
  42. # speed 三级合并(L1 conf < L2 ini [speed] < L3 CLI)
  43. tuning_conf_path = os.path.join(self.base_dir, 'conf', 'datax-tuning.conf')
  44. l1 = load_tuning_conf(tuning_conf_path)
  45. merged = merge_speed(l1, self.config_parser, self.cli_speed_overrides)
  46. # job.setting.speed 只放 channel:DataX 并发数按「全局 byte/record ÷ 每通道 byte/record」推导,
  47. # 且 byte/record 限制优先级高于 channel。若这里也放 byte/record(与 core 每通道同值),
  48. # 会算出 全局÷每通道=1 → 把 channel 压成 1 通道。byte/record 只作每通道限速放进 core.transport。
  49. speed = {
  50. JOB_SETTING_SPEED_CHANNEL: merged['channel'],
  51. }
  52. core_speed = {
  53. 'transport': {
  54. 'channel': {
  55. 'speed': {
  56. 'byte': merged['byte'],
  57. 'record': merged['record'],
  58. }
  59. }
  60. }
  61. }
  62. job_config_json = {
  63. 'job': {
  64. 'content': [
  65. {
  66. 'reader': self.get_reader(),
  67. 'writer': self.get_writer()
  68. }
  69. ],
  70. 'setting': {
  71. 'speed': speed
  72. }
  73. },
  74. 'core': core_speed
  75. }
  76. return job_config_json
  77. def run(self):
  78. job_config_json = self.assemble()
  79. # HDFS Mount的覆盖写入貌似有问题
  80. delete_file(self.output)
  81. with open(self.output, 'w') as w:
  82. json.dump(job_config_json, w, ensure_ascii=False)