This commit is contained in:
Maxime Beauchemin 2015-07-23 22:44:32 -07:00
Родитель 9e6ee42772
Коммит 1bfbd752f4
5 изменённых файлов: 13 добавлений и 11 удалений

Просмотреть файл

@ -26,14 +26,16 @@ class DbApiHook(BaseHook):
"""
supports_autocommit = False
def __init__(self, **kwargs):
try:
self.conn_id_name = kwargs[self.conn_name_attr]
except NameError:
def __init__(self, *args, **kwargs):
if not self.conn_name_attr:
raise AirflowException("conn_name_attr is not defined")
except KeyError:
raise AirflowException(
self.conn_name_attr + " was not passed in the kwargs")
elif len(args) == 1:
setattr(self, self.conn_name_attr, args[0])
elif self.conn_name_attr not in kwargs:
setattr(self, self.conn_name_attr, self.default_conn_name)
else:
setattr(self, self.conn_name_attr, kwargs[self.conn_name_attr])
def get_pandas_df(self, sql, parameters=None):
'''

Просмотреть файл

@ -16,7 +16,7 @@ class MySqlHook(DbApiHook):
"""
Returns a mysql connection object
"""
conn = self.get_connection(self.conn_id_name)
conn = self.get_connection(self.mysql_conn_id)
conn = MySQLdb.connect(
conn.host,
conn.login,

Просмотреть файл

@ -13,7 +13,7 @@ class PostgresHook(DbApiHook):
supports_autocommit = True
def get_conn(self):
conn = self.get_connection(self.conn_id_name)
conn = self.get_connection(self.postgres_conn_id)
return psycopg2.connect(
host=conn.host,
user=conn.login,

Просмотреть файл

@ -26,7 +26,7 @@ class PrestoHook(DbApiHook):
def get_conn(self):
"""Returns a connection object"""
db = self.get_connection(self.conn_id_name)
db = self.get_connection(self.presto_conn_id)
return presto.connect(
host=db.host,
port=db.port,

Просмотреть файл

@ -17,6 +17,6 @@ class SqliteHook(DbApiHook):
"""
Returns a sqlite connection object
"""
conn = self.get_connection(self.conn_id_name)
conn = self.get_connection(self.sqlite_conn_id)
conn = sqlite3.connect(conn.host)
return conn