postgresql_reader.py 4.7 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788
  1. # -*- coding:utf-8 -*-
  2. import re
  3. from configparser import ConfigParser
  4. from dw_base.datax.datax_constants import *
  5. from dw_base.datax.plugins.reader.reader import Reader
  6. # PostgreSQL reader
  7. from dw_base.datax.plugins.writer.postgresql_writer import POSTGRE_SQL_WRITER_PARAMETER_COLUMN, \
  8. POSTGRE_SQL_WRITER_PARAMETER_CONNECTION, POSTGRE_SQL_WRITER_PARAMETER_DATABASE, POSTGRE_SQL_WRITER_PARAMETER_TABLE
  9. POSTGRE_SQL_READER_NAME = 'postgresqlreader'
  10. POSTGRE_SQL_READER_PARAMETER_CONNECTION = 'connection'
  11. POSTGRE_SQL_READER_PARAMETER_DATABASE = 'database'
  12. POSTGRE_SQL_READER_PARAMETER_FETCH_SIZE = 'fetchSize'
  13. POSTGRE_SQL_READER_PARAMETER_QUERY_SQL = 'querySql'
  14. POSTGRE_SQL_READER_PARAMETER_TABLE = 'table'
  15. POSTGRE_SQL_READER_PARAMETER_COLUMN = 'column'
  16. POSTGRE_SQL_READER_PARAMETER_WHERE = 'where'
  17. POSTGRE_SQL_READER_PARAMETER_SPLIT_PK = 'splitPk'
  18. class PostgreSQLReader(Reader):
  19. def __init__(self, base_dir: str, config_parser: ConfigParser, start_date: str = None, stop_date: str = None):
  20. super(PostgreSQLReader, self).__init__(base_dir, config_parser, start_date, stop_date)
  21. self.plugin_name = POSTGRE_SQL_READER_NAME
  22. def load_others(self):
  23. start_date = self.start_date
  24. stop_date = self.stop_date
  25. database = self.config_parser.get(self.plugin_type, POSTGRE_SQL_WRITER_PARAMETER_DATABASE)
  26. self.check_config(POSTGRE_SQL_WRITER_PARAMETER_DATABASE, database)
  27. table = self.config_parser.get(self.plugin_type, POSTGRE_SQL_WRITER_PARAMETER_TABLE)
  28. self.check_config(POSTGRE_SQL_WRITER_PARAMETER_TABLE, table)
  29. fetch_size = self.config_parser.get(self.plugin_type, POSTGRE_SQL_READER_PARAMETER_FETCH_SIZE) or '1000'
  30. self.parameter[POSTGRE_SQL_READER_PARAMETER_FETCH_SIZE] = fetch_size
  31. split_pk = self.config_parser.get(self.plugin_type, POSTGRE_SQL_READER_PARAMETER_SPLIT_PK)
  32. self.parameter[POSTGRE_SQL_READER_PARAMETER_SPLIT_PK] = split_pk
  33. where = self.config_parser.get(self.plugin_type, POSTGRE_SQL_READER_PARAMETER_WHERE)
  34. where = where.replace('${start_date}', start_date)
  35. where = where.replace('${start-date}', start_date)
  36. where = where.replace('${stop_date}', stop_date)
  37. where = where.replace('${stop-date}', stop_date)
  38. self.parameter[POSTGRE_SQL_READER_PARAMETER_WHERE] = where
  39. jdbc_url: str = self.parameter[DS_POSTGRE_SQL_JDBC_URL]
  40. matcher = re.search('jdbc:postgresql://(.+?)/(.+)', jdbc_url)
  41. if matcher:
  42. if database:
  43. jdbc_url = jdbc_url.replace(matcher.group(2), database)
  44. elif jdbc_url.endswith('/'):
  45. jdbc_url = f'{jdbc_url}{database}'
  46. else:
  47. jdbc_url = f'{jdbc_url}/{database}'
  48. # 长连接抗断:writer 反压时 reader 空闲,aliyun rwlb 会掐空闲连接;默认无 socketTimeout
  49. # 时读操作会死等半开 socket(挂起不报错)。加 socketTimeout(秒)让它断了立即报错、
  50. # tcpKeepAlive 开保活探测。全局生效;日常小任务单批 fetch 很快,不受 300s 影响。
  51. _sep = '&' if '?' in jdbc_url else '?'
  52. jdbc_url = jdbc_url + _sep + 'tcpKeepAlive=true&socketTimeout=300'
  53. query_sql = self.config_parser.get(self.plugin_type, POSTGRE_SQL_READER_PARAMETER_QUERY_SQL)
  54. # 优先级:手写 querySql > [mask] 段自动生成 > table 透传
  55. if not query_sql and self.config_parser.has_section('mask'):
  56. from dw_base.datax.mask import build_query_sql
  57. columns_list = [c.strip() for c in self.config_parser.get(
  58. self.plugin_type, POSTGRE_SQL_WRITER_PARAMETER_COLUMN).split(',')]
  59. mask_config = dict(self.config_parser.items('mask'))
  60. query_sql = build_query_sql('postgresql', columns_list, mask_config, table, where)
  61. query_sql = query_sql.replace('${start_date}', start_date)
  62. query_sql = query_sql.replace('${start-date}', start_date)
  63. query_sql = query_sql.replace('${stop_date}', stop_date)
  64. query_sql = query_sql.replace('${stop-date}', stop_date)
  65. if query_sql:
  66. connection = {
  67. DS_POSTGRE_SQL_JDBC_URL: jdbc_url.split(','),
  68. POSTGRE_SQL_READER_PARAMETER_QUERY_SQL: query_sql.split(';')
  69. }
  70. else:
  71. connection = {
  72. DS_POSTGRE_SQL_JDBC_URL: jdbc_url.split(','),
  73. POSTGRE_SQL_READER_PARAMETER_TABLE: table.split(',')
  74. }
  75. self.parameter[POSTGRE_SQL_WRITER_PARAMETER_CONNECTION] = [connection]
  76. del self.parameter[DS_POSTGRE_SQL_JDBC_URL]
  77. def load_column(self):
  78. columns = self.config_parser.get(self.plugin_type, POSTGRE_SQL_WRITER_PARAMETER_COLUMN).split(',')
  79. self.parameter[POSTGRE_SQL_WRITER_PARAMETER_COLUMN] = columns