From 1814038721e22f7ea3d4cd6511448a0120cf7f27 Mon Sep 17 00:00:00 2001 From: Kesem Date: Wed, 17 Sep 2025 12:57:43 +0000 Subject: [PATCH 1/4] Replace sasl with pyhive.get_installed_sasl for pure-sasl compatibility - Remove direct sasl/saslwrapper imports - Use pyhive.hive.get_installed_sasl which handles pure-sasl compatibility - Maintain same SASL authentication logic for GSSAPI/Kerberos - Ensures compatibility with pure-sasl package while preserving functionality --- .../providers/apache/hive/hooks/hive.py | 41 ++++++++----------- 1 file changed, 16 insertions(+), 25 deletions(-) diff --git a/providers/apache/hive/src/airflow/providers/apache/hive/hooks/hive.py b/providers/apache/hive/src/airflow/providers/apache/hive/hooks/hive.py index 7ba8834694fff..33890906fa0e8 100644 --- a/providers/apache/hive/src/airflow/providers/apache/hive/hooks/hive.py +++ b/providers/apache/hive/src/airflow/providers/apache/hive/hooks/hive.py @@ -16,17 +16,26 @@ # specific language governing permissions and limitations # under the License. from __future__ import annotations - +from collections.abc import Iterable, Mapping import contextlib +import csv import os import re import socket import subprocess -import time -from collections.abc import Iterable, Mapping from tempfile import NamedTemporaryFile, TemporaryDirectory -from typing import TYPE_CHECKING, Any, Literal +import time +from typing import Any, Literal, TYPE_CHECKING +from airflow.configuration import conf +from airflow.exceptions import AirflowException, AirflowProviderDeprecationWarning +from airflow.providers.apache.hive.version_compat import ( + AIRFLOW_VAR_NAME_FORMAT_MAPPING, + BaseHook, +) +from airflow.providers.common.sql.hooks.sql import DbApiHook +from airflow.security import utils +from airflow.utils.helpers import as_flattened_list from deprecated import deprecated from typing_extensions import overload @@ -34,14 +43,7 @@ import pandas as pd import polars as pl -import csv -from airflow.configuration import conf -from airflow.exceptions import AirflowException, AirflowProviderDeprecationWarning -from airflow.providers.apache.hive.version_compat import AIRFLOW_VAR_NAME_FORMAT_MAPPING, BaseHook -from airflow.providers.common.sql.hooks.sql import DbApiHook -from airflow.security import utils -from airflow.utils.helpers import as_flattened_list HIVE_QUEUE_PRIORITIES = ["VERY_HIGH", "HIGH", "NORMAL", "LOW", "VERY_LOW"] @@ -573,21 +575,10 @@ def get_metastore_client(self) -> Any: conn_socket = TSocket.TSocket(host, conn.port) if conf.get("core", "security") == "kerberos" and auth_mechanism == "GSSAPI": - try: - import saslwrapper as sasl - except ImportError: - import sasl - - def sasl_factory() -> sasl.Client: - sasl_client = sasl.Client() - sasl_client.setAttr("host", host) - sasl_client.setAttr("service", kerberos_service_name) - sasl_client.init() - return sasl_client - from thrift_sasl import TSaslClientTransport - - transport = TSaslClientTransport(sasl_factory, "GSSAPI", conn_socket) + from pyhive.hive import get_installed_sasl + sasl_auth = 'GSSAPI' + transport = TSaslClientTransport(lambda: get_installed_sasl(host=host, sasl_auth=sasl_auth, service=kerberos_service_name), sasl_auth, conn_socket) else: transport = TTransport.TBufferedTransport(conn_socket) From 7436c0a5f90b831d783bd9e42047f31bad310c84 Mon Sep 17 00:00:00 2001 From: Kesem Date: Thu, 25 Sep 2025 12:54:58 +0000 Subject: [PATCH 2/4] fix ruff formating --- dags/simple_spark_xcom_test.py | 98 +++++++++++ dags/test_spark_kubernetes_xcom.py | 160 ++++++++++++++++++ .../providers/apache/hive/hooks/hive.py | 10 +- 3 files changed, 265 insertions(+), 3 deletions(-) create mode 100644 dags/simple_spark_xcom_test.py create mode 100644 dags/test_spark_kubernetes_xcom.py diff --git a/dags/simple_spark_xcom_test.py b/dags/simple_spark_xcom_test.py new file mode 100644 index 0000000000000..785bc11ab0cf4 --- /dev/null +++ b/dags/simple_spark_xcom_test.py @@ -0,0 +1,98 @@ +""" +Simple test DAG for SparkKubernetesOperator XCom functionality. + +This is a minimal test to verify that xcom push works with SparkKubernetesOperator. +""" + +from datetime import datetime, timedelta +from airflow import DAG +from airflow.providers.cncf.kubernetes.operators.spark_kubernetes import SparkKubernetesOperator +from airflow.operators.python import PythonOperator + +# Default arguments +default_args = { + "owner": "airflow", + "depends_on_past": False, + "start_date": datetime(2024, 1, 1), + "email_on_failure": False, + "email_on_retry": False, + "retries": 1, + "retry_delay": timedelta(minutes=5), +} + +# Create the DAG +dag = DAG( + "simple_spark_xcom_test", + default_args=default_args, + description="Simple test for SparkKubernetesOperator XCom", + schedule_interval=None, + catchup=False, + tags=["test", "spark", "xcom"], +) + +# Simple Spark application that just prints and returns data +simple_spark_template = { + "apiVersion": "sparkoperator.k8s.io/v1beta2", + "kind": "SparkApplication", + "metadata": {"name": "simple-spark-xcom-test", "namespace": "default"}, + "spec": { + "type": "Python", + "pythonVersion": "3", + "mode": "cluster", + "image": "apache/spark:3.5.0-python3", + "imagePullPolicy": "Always", + "mainApplicationFile": "local:///opt/spark/examples/src/main/python/pi.py", + "sparkVersion": "3.5.0", + "restartPolicy": {"type": "Never"}, + "driver": { + "cores": 1, + "coreLimit": "1200m", + "memory": "512m", + "labels": {"version": "3.5.0"}, + "serviceAccount": "spark", + }, + "executor": {"cores": 1, "instances": 1, "memory": "512m", "labels": {"version": "3.5.0"}}, + }, +} + +# Test task with xcom push enabled +test_spark_xcom = SparkKubernetesOperator( + task_id="test_spark_xcom", + name="simple-spark-xcom-test", + namespace="default", + template_spec=simple_spark_template, + do_xcom_push=True, # This is the key parameter to test! + get_logs=True, + delete_on_termination=True, + dag=dag, +) + + +# Task to check the xcom result +def check_xcom_result(**context): + """Check if xcom was successfully pushed.""" + print("=== Checking XCom Result ===") + + # Try to pull xcom from the Spark task + xcom_result = context["ti"].xcom_pull(task_ids="test_spark_xcom") + + print(f"XCom result type: {type(xcom_result)}") + print(f"XCom result: {xcom_result}") + + if xcom_result: + print("✅ SUCCESS: XCom was pushed successfully!") + print(f"✅ XCom contains: {xcom_result}") + else: + print("❌ FAILURE: No XCom data found") + + return xcom_result + + +check_result = PythonOperator( + task_id="check_xcom_result", + python_callable=check_xcom_result, + dag=dag, +) + +# Set up task dependencies +test_spark_xcom >> check_result diff --git a/dags/test_spark_kubernetes_xcom.py b/dags/test_spark_kubernetes_xcom.py new file mode 100644 index 0000000000000..75e6c0edd205a --- /dev/null +++ b/dags/test_spark_kubernetes_xcom.py @@ -0,0 +1,160 @@ +""" +Test DAG for SparkKubernetesOperator with XCom functionality. + +This DAG demonstrates the xcom push feature in SparkKubernetesOperator. +It creates a simple Spark job that returns data via xcom. +""" + +from datetime import datetime, timedelta +from airflow import DAG +from airflow.providers.cncf.kubernetes.operators.spark_kubernetes import SparkKubernetesOperator + +# Default arguments for the DAG +default_args = { + "owner": "airflow", + "depends_on_past": False, + "start_date": datetime(2024, 1, 1), + "email_on_failure": False, + "email_on_retry": False, + "retries": 1, + "retry_delay": timedelta(minutes=5), +} + +# Create the DAG +dag = DAG( + "test_spark_kubernetes_xcom", + default_args=default_args, + description="Test SparkKubernetesOperator with XCom functionality", + schedule_interval=None, # Manual trigger only + catchup=False, + tags=["test", "spark", "kubernetes", "xcom"], +) + +# Spark application template that will return data via xcom +spark_template = { + "apiVersion": "sparkoperator.k8s.io/v1beta2", + "kind": "SparkApplication", + "metadata": {"name": "test-spark-xcom-job", "namespace": "default"}, + "spec": { + "type": "Python", + "pythonVersion": "3", + "mode": "cluster", + "image": "apache/spark:3.5.0-python3", + "imagePullPolicy": "Always", + "mainApplicationFile": "local:///opt/spark/examples/src/main/python/pi.py", + "sparkVersion": "3.5.0", + "restartPolicy": {"type": "Never"}, + "driver": { + "cores": 1, + "coreLimit": "1200m", + "memory": "512m", + "labels": {"version": "3.5.0"}, + "serviceAccount": "spark", + }, + "executor": {"cores": 1, "instances": 1, "memory": "512m", "labels": {"version": "3.5.0"}}, + }, +} + +# Create a simple Python script that will be executed and return data +python_script = """ +import json +import sys +import os + +# Create some test data to return via xcom +test_data = { + "job_id": "test-spark-xcom-job", + "status": "completed", + "result": 3.14159, + "message": "Hello from Spark Kubernetes Operator!", + "timestamp": "2024-01-01T00:00:00Z" +} + +# Write the result to the xcom file +xcom_file = "/airflow/xcom/return.json" +os.makedirs(os.path.dirname(xcom_file), exist_ok=True) + +with open(xcom_file, 'w') as f: + json.dump(test_data, f) + +print("XCom data written successfully!") +print(f"Data: {test_data}") +""" + +# Create a custom Spark application that runs our Python script +custom_spark_template = { + "apiVersion": "sparkoperator.k8s.io/v1beta2", + "kind": "SparkApplication", + "metadata": {"name": "test-spark-xcom-custom-job", "namespace": "default"}, + "spec": { + "type": "Python", + "pythonVersion": "3", + "mode": "cluster", + "image": "apache/spark:3.5.0-python3", + "imagePullPolicy": "Always", + "mainApplicationFile": "local:///tmp/test_script.py", + "sparkVersion": "3.5.0", + "restartPolicy": {"type": "Never"}, + "driver": { + "cores": 1, + "coreLimit": "1200m", + "memory": "512m", + "labels": {"version": "3.5.0"}, + "serviceAccount": "spark", + }, + "executor": {"cores": 1, "instances": 1, "memory": "512m", "labels": {"version": "3.5.0"}}, + }, +} + +# Task 1: Test with built-in Spark example (Pi calculation) +test_spark_pi = SparkKubernetesOperator( + task_id="test_spark_pi_xcom", + name="test-spark-pi-xcom", + namespace="default", + template_spec=spark_template, + do_xcom_push=True, # Enable xcom push + get_logs=True, + delete_on_termination=True, + dag=dag, +) + +# Task 2: Test with custom Python script that returns structured data +test_spark_custom = SparkKubernetesOperator( + task_id="test_spark_custom_xcom", + name="test-spark-custom-xcom", + namespace="default", + template_spec=custom_spark_template, + do_xcom_push=True, # Enable xcom push + get_logs=True, + delete_on_termination=True, + dag=dag, +) + + +# Task 3: Print the xcom results +def print_xcom_results(**context): + """Print the xcom results from the Spark tasks.""" + print("=== XCom Results ===") + + # Get xcom from the Pi task + pi_result = context["ti"].xcom_pull(task_ids="test_spark_pi_xcom") + print(f"Pi task xcom result: {pi_result}") + + # Get xcom from the custom task + custom_result = context["ti"].xcom_pull(task_ids="test_spark_custom_xcom") + print(f"Custom task xcom result: {custom_result}") + + return {"pi_result": pi_result, "custom_result": custom_result} + + +from airflow.operators.python import PythonOperator + +print_results = PythonOperator( + task_id="print_xcom_results", + python_callable=print_xcom_results, + dag=dag, +) + +# Set up task dependencies +test_spark_pi >> print_results +test_spark_custom >> print_results diff --git a/providers/apache/hive/src/airflow/providers/apache/hive/hooks/hive.py b/providers/apache/hive/src/airflow/providers/apache/hive/hooks/hive.py index 33890906fa0e8..26af0d0b054c6 100644 --- a/providers/apache/hive/src/airflow/providers/apache/hive/hooks/hive.py +++ b/providers/apache/hive/src/airflow/providers/apache/hive/hooks/hive.py @@ -44,7 +44,6 @@ import polars as pl - HIVE_QUEUE_PRIORITIES = ["VERY_HIGH", "HIGH", "NORMAL", "LOW", "VERY_LOW"] @@ -577,8 +576,13 @@ def get_metastore_client(self) -> Any: if conf.get("core", "security") == "kerberos" and auth_mechanism == "GSSAPI": from thrift_sasl import TSaslClientTransport from pyhive.hive import get_installed_sasl - sasl_auth = 'GSSAPI' - transport = TSaslClientTransport(lambda: get_installed_sasl(host=host, sasl_auth=sasl_auth, service=kerberos_service_name), sasl_auth, conn_socket) + + sasl_auth = "GSSAPI" + transport = TSaslClientTransport( + lambda: get_installed_sasl(host=host, sasl_auth=sasl_auth, service=kerberos_service_name), + sasl_auth, + conn_socket, + ) else: transport = TTransport.TBufferedTransport(conn_socket) From 938708cb12fe242f57f2bad3c4bf02c098d1faf9 Mon Sep 17 00:00:00 2001 From: Kesem Date: Thu, 25 Sep 2025 16:28:50 +0300 Subject: [PATCH 3/4] cleanup: remove test DAG files to match UI version - Remove dags/simple_spark_xcom_test.py - Remove dags/test_spark_kubernetes_xcom.py - Keep only the Hive provider changes for pure-sasl compatibility --- dags/simple_spark_xcom_test.py | 98 ------------------ dags/test_spark_kubernetes_xcom.py | 160 ----------------------------- 2 files changed, 258 deletions(-) delete mode 100644 dags/simple_spark_xcom_test.py delete mode 100644 dags/test_spark_kubernetes_xcom.py diff --git a/dags/simple_spark_xcom_test.py b/dags/simple_spark_xcom_test.py deleted file mode 100644 index 785bc11ab0cf4..0000000000000 --- a/dags/simple_spark_xcom_test.py +++ /dev/null @@ -1,98 +0,0 @@ -""" -Simple test DAG for SparkKubernetesOperator XCom functionality. - -This is a minimal test to verify that xcom push works with SparkKubernetesOperator. -""" - -from datetime import datetime, timedelta -from airflow import DAG -from airflow.providers.cncf.kubernetes.operators.spark_kubernetes import SparkKubernetesOperator -from airflow.operators.python import PythonOperator - -# Default arguments -default_args = { - "owner": "airflow", - "depends_on_past": False, - "start_date": datetime(2024, 1, 1), - "email_on_failure": False, - "email_on_retry": False, - "retries": 1, - "retry_delay": timedelta(minutes=5), -} - -# Create the DAG -dag = DAG( - "simple_spark_xcom_test", - default_args=default_args, - description="Simple test for SparkKubernetesOperator XCom", - schedule_interval=None, - catchup=False, - tags=["test", "spark", "xcom"], -) - -# Simple Spark application that just prints and returns data -simple_spark_template = { - "apiVersion": "sparkoperator.k8s.io/v1beta2", - "kind": "SparkApplication", - "metadata": {"name": "simple-spark-xcom-test", "namespace": "default"}, - "spec": { - "type": "Python", - "pythonVersion": "3", - "mode": "cluster", - "image": "apache/spark:3.5.0-python3", - "imagePullPolicy": "Always", - "mainApplicationFile": "local:///opt/spark/examples/src/main/python/pi.py", - "sparkVersion": "3.5.0", - "restartPolicy": {"type": "Never"}, - "driver": { - "cores": 1, - "coreLimit": "1200m", - "memory": "512m", - "labels": {"version": "3.5.0"}, - "serviceAccount": "spark", - }, - "executor": {"cores": 1, "instances": 1, "memory": "512m", "labels": {"version": "3.5.0"}}, - }, -} - -# Test task with xcom push enabled -test_spark_xcom = SparkKubernetesOperator( - task_id="test_spark_xcom", - name="simple-spark-xcom-test", - namespace="default", - template_spec=simple_spark_template, - do_xcom_push=True, # This is the key parameter to test! - get_logs=True, - delete_on_termination=True, - dag=dag, -) - - -# Task to check the xcom result -def check_xcom_result(**context): - """Check if xcom was successfully pushed.""" - print("=== Checking XCom Result ===") - - # Try to pull xcom from the Spark task - xcom_result = context["ti"].xcom_pull(task_ids="test_spark_xcom") - - print(f"XCom result type: {type(xcom_result)}") - print(f"XCom result: {xcom_result}") - - if xcom_result: - print("✅ SUCCESS: XCom was pushed successfully!") - print(f"✅ XCom contains: {xcom_result}") - else: - print("❌ FAILURE: No XCom data found") - - return xcom_result - - -check_result = PythonOperator( - task_id="check_xcom_result", - python_callable=check_xcom_result, - dag=dag, -) - -# Set up task dependencies -test_spark_xcom >> check_result diff --git a/dags/test_spark_kubernetes_xcom.py b/dags/test_spark_kubernetes_xcom.py deleted file mode 100644 index 75e6c0edd205a..0000000000000 --- a/dags/test_spark_kubernetes_xcom.py +++ /dev/null @@ -1,160 +0,0 @@ -""" -Test DAG for SparkKubernetesOperator with XCom functionality. - -This DAG demonstrates the xcom push feature in SparkKubernetesOperator. -It creates a simple Spark job that returns data via xcom. -""" - -from datetime import datetime, timedelta -from airflow import DAG -from airflow.providers.cncf.kubernetes.operators.spark_kubernetes import SparkKubernetesOperator - -# Default arguments for the DAG -default_args = { - "owner": "airflow", - "depends_on_past": False, - "start_date": datetime(2024, 1, 1), - "email_on_failure": False, - "email_on_retry": False, - "retries": 1, - "retry_delay": timedelta(minutes=5), -} - -# Create the DAG -dag = DAG( - "test_spark_kubernetes_xcom", - default_args=default_args, - description="Test SparkKubernetesOperator with XCom functionality", - schedule_interval=None, # Manual trigger only - catchup=False, - tags=["test", "spark", "kubernetes", "xcom"], -) - -# Spark application template that will return data via xcom -spark_template = { - "apiVersion": "sparkoperator.k8s.io/v1beta2", - "kind": "SparkApplication", - "metadata": {"name": "test-spark-xcom-job", "namespace": "default"}, - "spec": { - "type": "Python", - "pythonVersion": "3", - "mode": "cluster", - "image": "apache/spark:3.5.0-python3", - "imagePullPolicy": "Always", - "mainApplicationFile": "local:///opt/spark/examples/src/main/python/pi.py", - "sparkVersion": "3.5.0", - "restartPolicy": {"type": "Never"}, - "driver": { - "cores": 1, - "coreLimit": "1200m", - "memory": "512m", - "labels": {"version": "3.5.0"}, - "serviceAccount": "spark", - }, - "executor": {"cores": 1, "instances": 1, "memory": "512m", "labels": {"version": "3.5.0"}}, - }, -} - -# Create a simple Python script that will be executed and return data -python_script = """ -import json -import sys -import os - -# Create some test data to return via xcom -test_data = { - "job_id": "test-spark-xcom-job", - "status": "completed", - "result": 3.14159, - "message": "Hello from Spark Kubernetes Operator!", - "timestamp": "2024-01-01T00:00:00Z" -} - -# Write the result to the xcom file -xcom_file = "/airflow/xcom/return.json" -os.makedirs(os.path.dirname(xcom_file), exist_ok=True) - -with open(xcom_file, 'w') as f: - json.dump(test_data, f) - -print("XCom data written successfully!") -print(f"Data: {test_data}") -""" - -# Create a custom Spark application that runs our Python script -custom_spark_template = { - "apiVersion": "sparkoperator.k8s.io/v1beta2", - "kind": "SparkApplication", - "metadata": {"name": "test-spark-xcom-custom-job", "namespace": "default"}, - "spec": { - "type": "Python", - "pythonVersion": "3", - "mode": "cluster", - "image": "apache/spark:3.5.0-python3", - "imagePullPolicy": "Always", - "mainApplicationFile": "local:///tmp/test_script.py", - "sparkVersion": "3.5.0", - "restartPolicy": {"type": "Never"}, - "driver": { - "cores": 1, - "coreLimit": "1200m", - "memory": "512m", - "labels": {"version": "3.5.0"}, - "serviceAccount": "spark", - }, - "executor": {"cores": 1, "instances": 1, "memory": "512m", "labels": {"version": "3.5.0"}}, - }, -} - -# Task 1: Test with built-in Spark example (Pi calculation) -test_spark_pi = SparkKubernetesOperator( - task_id="test_spark_pi_xcom", - name="test-spark-pi-xcom", - namespace="default", - template_spec=spark_template, - do_xcom_push=True, # Enable xcom push - get_logs=True, - delete_on_termination=True, - dag=dag, -) - -# Task 2: Test with custom Python script that returns structured data -test_spark_custom = SparkKubernetesOperator( - task_id="test_spark_custom_xcom", - name="test-spark-custom-xcom", - namespace="default", - template_spec=custom_spark_template, - do_xcom_push=True, # Enable xcom push - get_logs=True, - delete_on_termination=True, - dag=dag, -) - - -# Task 3: Print the xcom results -def print_xcom_results(**context): - """Print the xcom results from the Spark tasks.""" - print("=== XCom Results ===") - - # Get xcom from the Pi task - pi_result = context["ti"].xcom_pull(task_ids="test_spark_pi_xcom") - print(f"Pi task xcom result: {pi_result}") - - # Get xcom from the custom task - custom_result = context["ti"].xcom_pull(task_ids="test_spark_custom_xcom") - print(f"Custom task xcom result: {custom_result}") - - return {"pi_result": pi_result, "custom_result": custom_result} - - -from airflow.operators.python import PythonOperator - -print_results = PythonOperator( - task_id="print_xcom_results", - python_callable=print_xcom_results, - dag=dag, -) - -# Set up task dependencies -test_spark_pi >> print_results -test_spark_custom >> print_results From 5f1f6826aaea6df089657723b958993d2bae093b Mon Sep 17 00:00:00 2001 From: Kesem Date: Sun, 28 Sep 2025 10:15:33 +0000 Subject: [PATCH 4/4] Apply ruff formatting to hive.py - Reorganize imports according to ruff standards - Fix import ordering and formatting - Ensure code follows ruff linting rules - Auto-fix 2 ruff errors found in the file --- .../airflow/providers/apache/hive/hooks/hive.py | 14 ++++++++------ 1 file changed, 8 insertions(+), 6 deletions(-) diff --git a/providers/apache/hive/src/airflow/providers/apache/hive/hooks/hive.py b/providers/apache/hive/src/airflow/providers/apache/hive/hooks/hive.py index 26af0d0b054c6..4c8275d36cd0c 100644 --- a/providers/apache/hive/src/airflow/providers/apache/hive/hooks/hive.py +++ b/providers/apache/hive/src/airflow/providers/apache/hive/hooks/hive.py @@ -16,16 +16,20 @@ # specific language governing permissions and limitations # under the License. from __future__ import annotations -from collections.abc import Iterable, Mapping + import contextlib import csv import os import re import socket import subprocess -from tempfile import NamedTemporaryFile, TemporaryDirectory import time -from typing import Any, Literal, TYPE_CHECKING +from collections.abc import Iterable, Mapping +from tempfile import NamedTemporaryFile, TemporaryDirectory +from typing import TYPE_CHECKING, Any, Literal + +from deprecated import deprecated +from typing_extensions import overload from airflow.configuration import conf from airflow.exceptions import AirflowException, AirflowProviderDeprecationWarning @@ -36,8 +40,6 @@ from airflow.providers.common.sql.hooks.sql import DbApiHook from airflow.security import utils from airflow.utils.helpers import as_flattened_list -from deprecated import deprecated -from typing_extensions import overload if TYPE_CHECKING: import pandas as pd @@ -574,8 +576,8 @@ def get_metastore_client(self) -> Any: conn_socket = TSocket.TSocket(host, conn.port) if conf.get("core", "security") == "kerberos" and auth_mechanism == "GSSAPI": - from thrift_sasl import TSaslClientTransport from pyhive.hive import get_installed_sasl + from thrift_sasl import TSaslClientTransport sasl_auth = "GSSAPI" transport = TSaslClientTransport(