|
@@ -229,7 +229,7 @@ class HS2Api(Api):
|
|
|
|
|
|
|
|
@query_error_handler
|
|
@query_error_handler
|
|
|
def execute(self, notebook, snippet):
|
|
def execute(self, notebook, snippet):
|
|
|
- db = self._get_db(snippet)
|
|
|
|
|
|
|
+ db = self._get_db(snippet, cluster=snippet.get('selectedCompute'))
|
|
|
|
|
|
|
|
statement = self._get_current_statement(db, snippet)
|
|
statement = self._get_current_statement(db, snippet)
|
|
|
session = self._get_session(notebook, snippet['type'])
|
|
session = self._get_session(notebook, snippet['type'])
|
|
@@ -263,7 +263,7 @@ class HS2Api(Api):
|
|
|
@query_error_handler
|
|
@query_error_handler
|
|
|
def check_status(self, notebook, snippet):
|
|
def check_status(self, notebook, snippet):
|
|
|
response = {}
|
|
response = {}
|
|
|
- db = self._get_db(snippet)
|
|
|
|
|
|
|
+ db = self._get_db(snippet, cluster=snippet.get('selectedCompute'))
|
|
|
|
|
|
|
|
handle = self._get_handle(snippet)
|
|
handle = self._get_handle(snippet)
|
|
|
operation = db.get_operation_status(handle)
|
|
operation = db.get_operation_status(handle)
|
|
@@ -284,7 +284,7 @@ class HS2Api(Api):
|
|
|
|
|
|
|
|
@query_error_handler
|
|
@query_error_handler
|
|
|
def fetch_result(self, notebook, snippet, rows, start_over):
|
|
def fetch_result(self, notebook, snippet, rows, start_over):
|
|
|
- db = self._get_db(snippet)
|
|
|
|
|
|
|
+ db = self._get_db(snippet, cluster=snippet.get('selectedCompute'))
|
|
|
|
|
|
|
|
handle = self._get_handle(snippet)
|
|
handle = self._get_handle(snippet)
|
|
|
try:
|
|
try:
|
|
@@ -332,7 +332,7 @@ class HS2Api(Api):
|
|
|
|
|
|
|
|
@query_error_handler
|
|
@query_error_handler
|
|
|
def cancel(self, notebook, snippet):
|
|
def cancel(self, notebook, snippet):
|
|
|
- db = self._get_db(snippet)
|
|
|
|
|
|
|
+ db = self._get_db(snippet, cluster=snippet.get('selectedCompute'))
|
|
|
|
|
|
|
|
handle = self._get_handle(snippet)
|
|
handle = self._get_handle(snippet)
|
|
|
db.cancel_operation(handle)
|
|
db.cancel_operation(handle)
|
|
@@ -341,7 +341,7 @@ class HS2Api(Api):
|
|
|
|
|
|
|
|
@query_error_handler
|
|
@query_error_handler
|
|
|
def get_log(self, notebook, snippet, startFrom=None, size=None):
|
|
def get_log(self, notebook, snippet, startFrom=None, size=None):
|
|
|
- db = self._get_db(snippet)
|
|
|
|
|
|
|
+ db = self._get_db(snippet, cluster=snippet.get('selectedCompute'))
|
|
|
|
|
|
|
|
handle = self._get_handle(snippet)
|
|
handle = self._get_handle(snippet)
|
|
|
return db.get_log(handle, start_over=startFrom == 0)
|
|
return db.get_log(handle, start_over=startFrom == 0)
|
|
@@ -353,7 +353,7 @@ class HS2Api(Api):
|
|
|
from impala import conf as impala_conf
|
|
from impala import conf as impala_conf
|
|
|
|
|
|
|
|
if (snippet['type'] == 'hive' and beeswax_conf.CLOSE_QUERIES.get()) or (snippet['type'] == 'impala' and impala_conf.CLOSE_QUERIES.get()):
|
|
if (snippet['type'] == 'hive' and beeswax_conf.CLOSE_QUERIES.get()) or (snippet['type'] == 'impala' and impala_conf.CLOSE_QUERIES.get()):
|
|
|
- db = self._get_db(snippet)
|
|
|
|
|
|
|
+ db = self._get_db(snippet, cluster=snippet.get('selectedCompute'))
|
|
|
|
|
|
|
|
try:
|
|
try:
|
|
|
handle = self._get_handle(snippet)
|
|
handle = self._get_handle(snippet)
|
|
@@ -371,7 +371,7 @@ class HS2Api(Api):
|
|
|
@query_error_handler
|
|
@query_error_handler
|
|
|
def download(self, notebook, snippet, format, user_agent=None):
|
|
def download(self, notebook, snippet, format, user_agent=None):
|
|
|
try:
|
|
try:
|
|
|
- db = self._get_db(snippet)
|
|
|
|
|
|
|
+ db = self._get_db(snippet, cluster=snippet.get('selectedCompute'))
|
|
|
handle = self._get_handle(snippet)
|
|
handle = self._get_handle(snippet)
|
|
|
# Test handle to verify if still valid
|
|
# Test handle to verify if still valid
|
|
|
db.fetch(handle, start_over=True, rows=1)
|
|
db.fetch(handle, start_over=True, rows=1)
|
|
@@ -471,7 +471,7 @@ class HS2Api(Api):
|
|
|
|
|
|
|
|
@query_error_handler
|
|
@query_error_handler
|
|
|
def explain(self, notebook, snippet):
|
|
def explain(self, notebook, snippet):
|
|
|
- db = self._get_db(snippet)
|
|
|
|
|
|
|
+ db = self._get_db(snippet, cluster=snippet.get('selectedCompute'))
|
|
|
response = self._get_current_statement(db, snippet)
|
|
response = self._get_current_statement(db, snippet)
|
|
|
session = self._get_session(notebook, snippet['type'])
|
|
session = self._get_session(notebook, snippet['type'])
|
|
|
|
|
|
|
@@ -493,7 +493,7 @@ class HS2Api(Api):
|
|
|
|
|
|
|
|
@query_error_handler
|
|
@query_error_handler
|
|
|
def export_data_as_hdfs_file(self, snippet, target_file, overwrite):
|
|
def export_data_as_hdfs_file(self, snippet, target_file, overwrite):
|
|
|
- db = self._get_db(snippet)
|
|
|
|
|
|
|
+ db = self._get_db(snippet, cluster=snippet.get('selectedCompute'))
|
|
|
|
|
|
|
|
handle = self._get_handle(snippet)
|
|
handle = self._get_handle(snippet)
|
|
|
max_rows = DOWNLOAD_ROW_LIMIT.get()
|
|
max_rows = DOWNLOAD_ROW_LIMIT.get()
|
|
@@ -505,7 +505,7 @@ class HS2Api(Api):
|
|
|
|
|
|
|
|
|
|
|
|
|
def export_data_as_table(self, notebook, snippet, destination, is_temporary=False, location=None):
|
|
def export_data_as_table(self, notebook, snippet, destination, is_temporary=False, location=None):
|
|
|
- db = self._get_db(snippet)
|
|
|
|
|
|
|
+ db = self._get_db(snippet, cluster=snippet.get('selectedCompute'))
|
|
|
|
|
|
|
|
response = self._get_current_statement(db, snippet)
|
|
response = self._get_current_statement(db, snippet)
|
|
|
session = self._get_session(notebook, snippet['type'])
|
|
session = self._get_session(notebook, snippet['type'])
|
|
@@ -529,7 +529,7 @@ class HS2Api(Api):
|
|
|
|
|
|
|
|
|
|
|
|
|
def export_large_data_to_hdfs(self, notebook, snippet, destination):
|
|
def export_large_data_to_hdfs(self, notebook, snippet, destination):
|
|
|
- db = self._get_db(snippet)
|
|
|
|
|
|
|
+ db = self._get_db(snippet, cluster=snippet.get('selectedCompute'))
|
|
|
|
|
|
|
|
response = self._get_current_statement(db, snippet)
|
|
response = self._get_current_statement(db, snippet)
|
|
|
session = self._get_session(notebook, snippet['type'])
|
|
session = self._get_session(notebook, snippet['type'])
|
|
@@ -563,7 +563,7 @@ DROP TABLE IF EXISTS `%(table)s`;
|
|
|
|
|
|
|
|
|
|
|
|
|
def statement_risk(self, notebook, snippet):
|
|
def statement_risk(self, notebook, snippet):
|
|
|
- db = self._get_db(snippet)
|
|
|
|
|
|
|
+ db = self._get_db(snippet, cluster=snippet.get('selectedCompute'))
|
|
|
|
|
|
|
|
response = self._get_current_statement(db, snippet)
|
|
response = self._get_current_statement(db, snippet)
|
|
|
query = response['statement']
|
|
query = response['statement']
|
|
@@ -574,7 +574,7 @@ DROP TABLE IF EXISTS `%(table)s`;
|
|
|
|
|
|
|
|
|
|
|
|
|
def statement_compatibility(self, notebook, snippet, source_platform, target_platform):
|
|
def statement_compatibility(self, notebook, snippet, source_platform, target_platform):
|
|
|
- db = self._get_db(snippet)
|
|
|
|
|
|
|
+ db = self._get_db(snippet, cluster=snippet.get('selectedCompute'))
|
|
|
|
|
|
|
|
response = self._get_current_statement(db, snippet)
|
|
response = self._get_current_statement(db, snippet)
|
|
|
query = response['statement']
|
|
query = response['statement']
|
|
@@ -585,7 +585,7 @@ DROP TABLE IF EXISTS `%(table)s`;
|
|
|
|
|
|
|
|
|
|
|
|
|
def statement_similarity(self, notebook, snippet, source_platform):
|
|
def statement_similarity(self, notebook, snippet, source_platform):
|
|
|
- db = self._get_db(snippet)
|
|
|
|
|
|
|
+ db = self._get_db(snippet, cluster=snippet.get('selectedCompute'))
|
|
|
|
|
|
|
|
response = self._get_current_statement(db, snippet)
|
|
response = self._get_current_statement(db, snippet)
|
|
|
query = response['statement']
|
|
query = response['statement']
|
|
@@ -736,11 +736,11 @@ DROP TABLE IF EXISTS `%(table)s`;
|
|
|
|
|
|
|
|
|
|
|
|
|
def get_browse_query(self, snippet, database, table, partition_spec=None):
|
|
def get_browse_query(self, snippet, database, table, partition_spec=None):
|
|
|
- db = self._get_db(snippet)
|
|
|
|
|
|
|
+ db = self._get_db(snippet, cluster=snippet.get('selectedCompute'))
|
|
|
table = db.get_table(database, table)
|
|
table = db.get_table(database, table)
|
|
|
if table.is_impala_only:
|
|
if table.is_impala_only:
|
|
|
snippet['type'] = 'impala'
|
|
snippet['type'] = 'impala'
|
|
|
- db = self._get_db(snippet)
|
|
|
|
|
|
|
+ db = self._get_db(snippet, cluster=snippet.get('selectedCompute'))
|
|
|
|
|
|
|
|
if partition_spec is not None:
|
|
if partition_spec is not None:
|
|
|
decoded_spec = urllib.unquote(partition_spec)
|
|
decoded_spec = urllib.unquote(partition_spec)
|