meta.py 10 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210
  1. # Licensed to the Apache Software Foundation (ASF) under one or more
  2. # contributor license agreements. See the NOTICE file distributed with
  3. # this work for additional information regarding copyright ownership.
  4. # The ASF licenses this file to You under the Apache License, Version 2.0
  5. # (the "License"); you may not use this file except in compliance with
  6. # the License. You may obtain a copy of the License at
  7. #
  8. # http://www.apache.org/licenses/LICENSE-2.0
  9. #
  10. # Unless required by applicable law or agreed to in writing, software
  11. # distributed under the License is distributed on an "AS IS" BASIS,
  12. # WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
  13. # See the License for the specific language governing permissions and
  14. # limitations under the License.
  15. import sys
  16. import logging
  17. from phoenixdb.avatica.proto import common_pb2
  18. from phoenixdb.errors import ProgrammingError
  19. from phoenixdb.cursor import DictCursor
  20. __all__ = ['Meta']
  21. logger = logging.getLogger(__name__)
  22. class Meta(object):
  23. """Database meta for querying MetaData
  24. """
  25. def __init__(self, connection):
  26. self._connection = connection
  27. def get_catalogs(self):
  28. if self._connection._closed:
  29. raise ProgrammingError('The connection is already closed.')
  30. result = self._connection._client.get_catalogs(self._connection._id)
  31. with DictCursor(self._connection) as cursor:
  32. cursor._process_result(result)
  33. return cursor.fetchall()
  34. def get_schemas(self, catalog=None, schemaPattern=None):
  35. if self._connection._closed:
  36. raise ProgrammingError('The connection is already closed.')
  37. result = self._connection._client.get_schemas(self._connection._id, catalog, schemaPattern)
  38. with DictCursor(self._connection) as cursor:
  39. cursor._process_result(result)
  40. return self._fix_default(cursor.fetchall(), schemaPattern=schemaPattern)
  41. def get_tables(self, catalog=None, schemaPattern=None, tableNamePattern=None, typeList=None):
  42. if self._connection._closed:
  43. raise ProgrammingError('The connection is already closed.')
  44. result = self._connection._client.get_tables(
  45. self._connection._id, catalog, schemaPattern, tableNamePattern, typeList=typeList)
  46. with DictCursor(self._connection) as cursor:
  47. cursor._process_result(result)
  48. return self._fix_default(cursor.fetchall(), catalog, schemaPattern)
  49. def get_columns(self, catalog=None, schemaPattern=None, tableNamePattern=None,
  50. columnNamePattern=None):
  51. if self._connection._closed:
  52. raise ProgrammingError('The connection is already closed.')
  53. result = self._connection._client.get_columns(
  54. self._connection._id, catalog, schemaPattern, tableNamePattern, columnNamePattern)
  55. with DictCursor(self._connection) as cursor:
  56. cursor._process_result(result)
  57. return self._fix_default(cursor.fetchall(), catalog, schemaPattern)
  58. def get_table_types(self):
  59. if self._connection._closed:
  60. raise ProgrammingError('The connection is already closed.')
  61. result = self._connection._client.get_table_types(self._connection._id)
  62. with DictCursor(self._connection) as cursor:
  63. cursor._process_result(result)
  64. return cursor.fetchall()
  65. def get_type_info(self):
  66. if self._connection._closed:
  67. raise ProgrammingError('The connection is already closed.')
  68. result = self._connection._client.get_type_info(self._connection._id)
  69. with DictCursor(self._connection) as cursor:
  70. cursor._process_result(result)
  71. return cursor.fetchall()
  72. def get_primary_keys(self, catalog=None, schema=None, table=None):
  73. if self._connection._closed:
  74. raise ProgrammingError('The cursor is already closed.')
  75. state = common_pb2.QueryState()
  76. state.type = common_pb2.StateType.METADATA
  77. state.op = common_pb2.MetaDataOperation.GET_PRIMARY_KEYS
  78. state.has_args = True
  79. state.has_op = True
  80. catalog_arg = self._moa_string_arg_factory(catalog)
  81. schema_arg = self._moa_string_arg_factory(schema)
  82. table_arg = self._moa_string_arg_factory(table)
  83. state.args.extend([catalog_arg, schema_arg, table_arg])
  84. with DictCursor(self._connection) as cursor:
  85. syncResultResponse = cursor.get_sync_results(state)
  86. if not syncResultResponse.more_results:
  87. return []
  88. signature = common_pb2.Signature()
  89. signature.columns.append(self._column_meta_data_factory(1, 'TABLE_CAT', 12))
  90. signature.columns.append(self._column_meta_data_factory(2, 'TABLE_SCHEM', 12))
  91. signature.columns.append(self._column_meta_data_factory(3, 'TABLE_NAME', 12))
  92. signature.columns.append(self._column_meta_data_factory(4, 'COLUMN_NAME', 12))
  93. signature.columns.append(self._column_meta_data_factory(5, 'KEY_SEQ', 5))
  94. signature.columns.append(self._column_meta_data_factory(6, 'PK_NAME', 12))
  95. # The following are non-standard Phoenix extensions
  96. # This returns '\x00\x00\x00A' or '\x00\x00\x00D' , but that's consistent with Java
  97. signature.columns.append(self._column_meta_data_factory(7, 'ASC_OR_DESC', 12))
  98. signature.columns.append(self._column_meta_data_factory(8, 'DATA_TYPE', 5))
  99. signature.columns.append(self._column_meta_data_factory(9, 'TYPE_NAME', 12))
  100. signature.columns.append(self._column_meta_data_factory(10, 'COLUMN_SIZE', 5))
  101. signature.columns.append(self._column_meta_data_factory(11, 'TYPE_ID', 5))
  102. signature.columns.append(self._column_meta_data_factory(12, 'VIEW_CONSTANT', 12))
  103. cursor.fetch(signature)
  104. return cursor.fetchall()
  105. def get_index_info(self, catalog=None, schema=None, table=None, unique=False, approximate=False):
  106. if self._connection._closed:
  107. raise ProgrammingError('The cursor is already closed.')
  108. state = common_pb2.QueryState()
  109. state.type = common_pb2.StateType.METADATA
  110. state.op = common_pb2.MetaDataOperation.GET_INDEX_INFO
  111. state.has_args = True
  112. state.has_op = True
  113. catalog_arg = self._moa_string_arg_factory(catalog)
  114. schema_arg = self._moa_string_arg_factory(schema)
  115. table_arg = self._moa_string_arg_factory(table)
  116. unique_arg = self._moa_bool_arg_factory(unique)
  117. approximate_arg = self._moa_bool_arg_factory(approximate)
  118. state.args.extend([catalog_arg, schema_arg, table_arg, unique_arg, approximate_arg])
  119. with DictCursor(self._connection) as cursor:
  120. syncResultResponse = cursor.get_sync_results(state)
  121. if not syncResultResponse.more_results:
  122. return []
  123. signature = common_pb2.Signature()
  124. signature.columns.append(self._column_meta_data_factory(1, 'TABLE_CAT', 12))
  125. signature.columns.append(self._column_meta_data_factory(2, 'TABLE_SCHEM', 12))
  126. signature.columns.append(self._column_meta_data_factory(3, 'TABLE_NAME', 12))
  127. signature.columns.append(self._column_meta_data_factory(4, 'NON_UNIQUE', 16))
  128. signature.columns.append(self._column_meta_data_factory(5, 'INDEX_QUALIFIER', 12))
  129. signature.columns.append(self._column_meta_data_factory(6, 'INDEX_NAME', 12))
  130. signature.columns.append(self._column_meta_data_factory(7, 'TYPE', 5))
  131. signature.columns.append(self._column_meta_data_factory(8, 'ORDINAL_POSITION', 5))
  132. signature.columns.append(self._column_meta_data_factory(9, 'COLUMN_NAME', 12))
  133. signature.columns.append(self._column_meta_data_factory(10, 'ASC_OR_DESC', 12))
  134. signature.columns.append(self._column_meta_data_factory(11, 'CARDINALITY', 5))
  135. signature.columns.append(self._column_meta_data_factory(12, 'PAGES', 5))
  136. signature.columns.append(self._column_meta_data_factory(13, 'FILTER_CONDITION', 12))
  137. # The following are non-standard Phoenix extensions
  138. signature.columns.append(self._column_meta_data_factory(14, 'DATA_TYPE', 5))
  139. signature.columns.append(self._column_meta_data_factory(15, 'TYPE_NAME', 12))
  140. signature.columns.append(self._column_meta_data_factory(16, 'TYPE_ID', 5))
  141. signature.columns.append(self._column_meta_data_factory(17, 'COLUMN_FAMILY', 12))
  142. signature.columns.append(self._column_meta_data_factory(18, 'COLUMN_SIZE', 5))
  143. signature.columns.append(self._column_meta_data_factory(19, 'ARRAY_SIZE', 5))
  144. cursor.fetch(signature)
  145. return cursor.fetchall()
  146. def _column_meta_data_factory(self, ordinal, column_name, jdbc_code):
  147. cmd = common_pb2.ColumnMetaData()
  148. cmd.ordinal = ordinal
  149. cmd.column_name = column_name
  150. cmd.type.id = jdbc_code
  151. cmd.nullable = 2
  152. return cmd
  153. def _moa_string_arg_factory(self, arg):
  154. moa = common_pb2.MetaDataOperationArgument()
  155. if arg is None:
  156. moa.type = common_pb2.MetaDataOperationArgument.ArgumentType.NULL
  157. else:
  158. moa.type = common_pb2.MetaDataOperationArgument.ArgumentType.STRING
  159. moa.string_value = arg
  160. return moa
  161. def _moa_bool_arg_factory(self, arg):
  162. moa = common_pb2.MetaDataOperationArgument()
  163. if arg is None:
  164. moa.type = common_pb2.MetaDataOperationArgument.ArgumentType.NULL
  165. else:
  166. moa.type = common_pb2.MetaDataOperationArgument.ArgumentType.BOOL
  167. moa.bool_value = arg
  168. return moa
  169. def _fix_default(self, rows, catalog=None, schemaPattern=None):
  170. '''Workaround for PHOENIX-6003'''
  171. if schemaPattern == '':
  172. rows = [row for row in rows if row['TABLE_SCHEM'] is None]
  173. if catalog == '':
  174. rows = [row for row in rows if row['TABLE_CATALOG'] is None]
  175. # Couldn't find a sane way to do it that works on 2 and 3
  176. if sys.version_info.major == 3:
  177. return [{k: v or '' for k, v in row.items()} for row in rows]
  178. else:
  179. return [{k: v or '' for k, v in row.iteritems()} for row in rows]