diff --git a/airflow/plugins/custom_operators/iceberg_operator.py b/airflow/plugins/custom_operators/iceberg_operator.py index 3dd0d45..7310a4a 100644 --- a/airflow/plugins/custom_operators/iceberg_operator.py +++ b/airflow/plugins/custom_operators/iceberg_operator.py @@ -58,9 +58,11 @@ def __init__( def _get_spark_conf(self) -> Dict[str, str]: """Get Spark configuration for Iceberg.""" return { - 'spark.jars.packages': 'org.apache.iceberg:iceberg-spark-runtime-4.0_2.13:1.10.0,' - 'org.apache.hadoop:hadoop-aws:3.4.1', + 'spark.jars.packages': 'org.apache.iceberg:iceberg-spark-runtime-3.5_2.12:1.10.0,' + 'org.apache.hadoop:hadoop-aws:3.3.4', 'spark.sql.extensions': 'org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions', + 'spark.sql.catalog.spark_catalog': 'org.apache.iceberg.spark.SparkSessionCatalog', + 'spark.sql.catalog.spark_catalog.type': 'hive', 'spark.sql.catalog.local': 'org.apache.iceberg.spark.SparkCatalog', 'spark.sql.catalog.local.type': 'hadoop', 'spark.sql.catalog.local.warehouse': self.warehouse_path, diff --git a/docker/airflow/Dockerfile b/docker/airflow/Dockerfile index ec6d1eb..9d0f3a0 100644 --- a/docker/airflow/Dockerfile +++ b/docker/airflow/Dockerfile @@ -1,4 +1,4 @@ -FROM apache/airflow:3.1.7-python3.13 +FROM apache/airflow:3.1.5-python3.13 USER root diff --git a/docker/airflow/requirements.txt b/docker/airflow/requirements.txt index e650be3..6e7b1aa 100644 --- a/docker/airflow/requirements.txt +++ b/docker/airflow/requirements.txt @@ -1,4 +1,4 @@ -# Airflow providers (Airflow 3.1 compatible) +# Airflow providers (Airflow 3.0 compatible) apache-airflow-providers-apache-spark>=5.0.0 apache-airflow-providers-amazon>=9.0.0 diff --git a/spark/config/log4j.properties b/spark/config/log4j.properties index e09ada9..a47aa88 100644 --- a/spark/config/log4j.properties +++ b/spark/config/log4j.properties @@ -15,3 +15,5 @@ log4j.logger.org.apache.spark.repl.SparkIMain$exprTyper=INFO log4j.logger.org.apache.spark.repl.SparkILoop$SparkILoopInterpreter=INFO log4j.logger.org.apache.parquet=ERROR log4j.logger.parquet=ERROR +log4j.logger.org.apache.hadoop.hive.metastore.RetryingHMSHandler=FATAL +log4j.logger.org.apache.hadoop.hive.ql.exec.FunctionRegistry=ERROR diff --git a/spark/config/spark-defaults.conf b/spark/config/spark-defaults.conf index 7b50065..47790b1 100644 --- a/spark/config/spark-defaults.conf +++ b/spark/config/spark-defaults.conf @@ -23,6 +23,8 @@ spark.hadoop.fs.s3a.connection.ssl.enabled false # Iceberg Configuration spark.sql.extensions org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions +spark.sql.catalog.spark_catalog org.apache.iceberg.spark.SparkSessionCatalog +spark.sql.catalog.spark_catalog.type hive spark.sql.catalog.local org.apache.iceberg.spark.SparkCatalog spark.sql.catalog.local.type hadoop spark.sql.catalog.local.warehouse s3a://warehouse/ diff --git a/terraform/modules/airflow/main.tf b/terraform/modules/airflow/main.tf index a2c5c89..e326a1d 100644 --- a/terraform/modules/airflow/main.tf +++ b/terraform/modules/airflow/main.tf @@ -3,16 +3,16 @@ resource "helm_release" "airflow" { repository = "https://airflow.apache.org" chart = "airflow" namespace = var.namespace - version = "1.19.0" + version = "1.18.0" values = [ <<-EOT executor: "KubernetesExecutor" - airflowVersion: "3.1.7" + airflowVersion: "3.0.2" defaultAirflowRepository: apache/airflow - defaultAirflowTag: "3.1.7" + defaultAirflowTag: "3.0.2" webserver: service: diff --git a/terraform/modules/spark/main.tf b/terraform/modules/spark/main.tf index ffdc9c7..960db49 100644 --- a/terraform/modules/spark/main.tf +++ b/terraform/modules/spark/main.tf @@ -33,7 +33,7 @@ resource "kubernetes_stateful_set_v1" "spark_master" { container { name = "spark-master" - image = "apache/spark:4.0.1" + image = "apache/spark:3.5.0" command = ["/opt/spark/bin/spark-class"] args = ["org.apache.spark.deploy.master.Master"] @@ -156,7 +156,7 @@ resource "kubernetes_deployment_v1" "spark_worker" { spec { container { name = "spark-worker" - image = "apache/spark:4.0.1" + image = "apache/spark:3.5.0" command = ["/opt/spark/bin/spark-class"] args = ["org.apache.spark.deploy.worker.Worker", "spark://spark-master:7077"]