diff --git a/.gitignore b/.gitignore index 4380808..43085e1 100644 --- a/.gitignore +++ b/.gitignore @@ -1,2 +1,3 @@ .idea/ -.conveyor \ No newline at end of file +.conveyor +.terraform/ \ No newline at end of file diff --git a/basic/pi_spark/Dockerfile b/basic/pi_spark/Dockerfile index e4cd490..8d6a4df 100644 --- a/basic/pi_spark/Dockerfile +++ b/basic/pi_spark/Dockerfile @@ -1,4 +1,4 @@ -FROM public.ecr.aws/dataminded/spark-k8s-glue:v3.5.1-hadoop-3.3.6-v1 +FROM public.ecr.aws/dataminded/spark-k8s-glue:v3.5.6-hadoop-3.3.6-v1 ENV PYSPARK_PYTHON python3 WORKDIR /opt/spark/work-dir diff --git a/basic/pi_spark/dags/pi_spark.py b/basic/pi_spark/dags/pi_spark.py index 2d9ccbe..12e4ba2 100644 --- a/basic/pi_spark/dags/pi_spark.py +++ b/basic/pi_spark/dags/pi_spark.py @@ -1,33 +1,33 @@ from airflow import DAG from conveyor.operators import ConveyorSparkSubmitOperatorV2 -from datetime import timedelta +from datetime import timedelta, datetime from airflow.utils import dates default_args = { "owner": "Conveyor", "depends_on_past": False, - "start_date": dates.days_ago(2), "email": [], "email_on_failure": False, "email_on_retry": False, "retries": 0, + "start_date": datetime(year=2026, month=1, day=5), "retry_delay": timedelta(minutes=5), } dag = DAG( "samples_pi_spark", default_args=default_args, - schedule_interval="@daily", + schedule="@daily", max_active_runs=1, ) role = "conveyor-samples" -sample_task = ConveyorSparkSubmitOperatorV2( +ConveyorSparkSubmitOperatorV2( dag=dag, task_id="calculate_pi", - num_executors="4", + num_executors=2, driver_instance_type="mx.medium", - executor_instance_type="mx.medium", + executor_instance_type="cx.xlarge", instance_life_cycle="spot", # Other options are `on-demand`, `driver-on-demand-executors-spot` aws_role=role, spark_main_version=3, @@ -35,6 +35,26 @@ application_args=[ "--date", "{{ ds }}", "--env", "{{ macros.conveyor.env() }}", - "--partitions", "2000", + "--partitions", "1000", + "--iterations", "3000000000", + ], +) + + +ConveyorSparkSubmitOperatorV2( + dag=dag, + task_id="calculate_pi_inefficient", + num_executors=1, + driver_instance_type="mx.medium", + executor_instance_type="cx.xlarge", + instance_life_cycle="spot", # Other options are `on-demand`, `driver-on-demand-executors-spot` + aws_role=role, + spark_main_version=3, + application="local:///opt/spark/work-dir/src/pi_spark/app.py", + application_args=[ + "--date", "{{ ds }}", + "--env", "{{ macros.conveyor.env() }}", + "--partitions", "3", + "--iterations", "3000000000", ], ) diff --git a/basic/pi_spark/dev-requirements.in b/basic/pi_spark/dev-requirements.in index 39073ac..d4f6524 100644 --- a/basic/pi_spark/dev-requirements.in +++ b/basic/pi_spark/dev-requirements.in @@ -1,5 +1,5 @@ -c requirements.txt -pyspark==3.5.1 +pyspark==3.5.6 pandas==1.5.2 pytest>=6.0 pytest-cov diff --git a/basic/pi_spark/dev-requirements.txt b/basic/pi_spark/dev-requirements.txt index e140221..3d91581 100644 --- a/basic/pi_spark/dev-requirements.txt +++ b/basic/pi_spark/dev-requirements.txt @@ -5,9 +5,9 @@ black==22.12.0 click==8.1.3 # via black coverage[toml]==7.0.5 - # via - # coverage - # pytest-cov + # via pytest-cov +exceptiongroup==1.3.1 + # via pytest flake8==6.0.0 # via -r dev-requirements.in iniconfig==2.0.0 @@ -34,7 +34,7 @@ pycodestyle==2.10.0 # via flake8 pyflakes==3.0.1 # via flake8 -pyspark==3.5.1 +pyspark==3.5.6 # via -r dev-requirements.in pytest==7.2.1 # via @@ -48,3 +48,10 @@ pytz==2022.7.1 # via pandas six==1.16.0 # via python-dateutil +tomli==2.3.0 + # via + # black + # coverage + # pytest +typing-extensions==4.15.0 + # via exceptiongroup diff --git a/basic/pi_spark/ide.yaml b/basic/pi_spark/ide.yaml index 5239c94..ac38092 100644 --- a/basic/pi_spark/ide.yaml +++ b/basic/pi_spark/ide.yaml @@ -3,21 +3,3 @@ buildSteps: cmd: | sudo mkdir -p /usr/share/man/man1 sudo apt install make docker-compose openjdk-11-jre gcc libbz2-dev openssl libncurses5-dev libncursesw5-dev libssl-dev libreadline-dev liblzma-dev libsqlite3-dev -y - - name: Install pyenv - cmd: | - curl https://pyenv.run | bash - - #Add to bashrc - echo 'export PYENV_ROOT="$HOME/.pyenv"' >> ~/.bashrc - echo 'command -v pyenv >/dev/null || export PATH="$PYENV_ROOT/bin:$PATH"' >> ~/.bashrc - echo 'eval "$(pyenv init -)"' >> ~/.bashrc - - #Add to profile - echo 'export PYENV_ROOT="$HOME/.pyenv"' >> ~/.profile - echo 'command -v pyenv >/dev/null || export PATH="$PYENV_ROOT/bin:$PATH"' >> ~/.profile - echo 'eval "$(pyenv init -)"' >> ~/.profile - - name: Install python 3.9.17 - cmd: | - . ~/.profile - pyenv install 3.9.17 - pyenv global 3.9.17 diff --git a/basic/pi_spark/src/pi_spark/app.py b/basic/pi_spark/src/pi_spark/app.py index f68c1ce..f1f252c 100644 --- a/basic/pi_spark/src/pi_spark/app.py +++ b/basic/pi_spark/src/pi_spark/app.py @@ -10,28 +10,52 @@ from pi_spark.common.spark import ClosableSparkSession, transform, SparkLogger - DataFrame.transform = transform def main(): parser = argparse.ArgumentParser(description="pi_spark") parser.add_argument( - "-d", "--date", dest="date", help="date in format YYYY-mm-dd", required=True + "-d", + "--date", + dest="date", + help="date in format YYYY-mm-dd", + required=True, + ) + parser.add_argument( + "-e", + "--env", + dest="env", + help="environment we are executing in", + required=True, ) parser.add_argument( - "-e", "--env", dest="env", help="environment we are executing in", required=True - ) + "-p", + "--partitions", + dest="partitions", + help="number of partitions to calculate pi", + required=False, + ) parser.add_argument( - "-p", "--partitions", dest="partitions", help="number of partitions to calculate pi", required=False + "-i", + "--iterations", + dest="iterations", + help="number of iterations per partition to calculate pi", + required=False, ) args = parser.parse_args() - with ClosableSparkSession("pi_spark") as session: - run(session, args.env, args.date, int(args.partitions)) + with ClosableSparkSession("pi_spark") as session: + run(session, args.env, args.date, int(args.partitions), int(args.iterations)) -def run(spark: SparkSession, environment: str, date: str, partitions: int=2): +def run( + spark: SparkSession, + environment: str, + date: str, + partitions: int = 1000, + iterations: int = 10000000000, +): """Main ETL script definition. :return: None @@ -39,18 +63,16 @@ def run(spark: SparkSession, environment: str, date: str, partitions: int=2): # execute ETL pipeline logger = SparkLogger(spark) logger.info(f"Executing job for {environment} on {date}") - n = 100000 * partitions logger.info(f"Partitions: {partitions}") - logger.info(f"n: {n}") - def f(_: int) -> float: - x = random() * 2 - 1 - y = random() * 2 - 1 - return 1 if x ** 2 + y ** 2 <= 1 else 0 + logger.info(f"number of iterations/samples: {iterations}") - count = spark.sparkContext.parallelize(range(1, n + 1), partitions).map(f).reduce(add) - logger.info("Pi is roughly %f" % (4.0 * count / n)) + def f(_: int) -> float: + x, y = random(), random() + return x * x + y * y < 1 + count = spark.sparkContext.parallelize(range(0, iterations), partitions).filter(f).count() + logger.info("Pi is roughly %f" % (4.0 * count / iterations)) if __name__ == "__main__": - main() \ No newline at end of file + main() diff --git a/basic/pi_spark/src/pi_spark/common/spark.py b/basic/pi_spark/src/pi_spark/common/spark.py index 36ff55d..ea2c112 100644 --- a/basic/pi_spark/src/pi_spark/common/spark.py +++ b/basic/pi_spark/src/pi_spark/common/spark.py @@ -72,7 +72,6 @@ def __enter__(self): spark_builder.config("spark.sql.hive.metastorePartitionPruning", "false") spark_builder.config("spark.sql.hive.convertMetastoreParquet", "false") - # add other config params for key, val in self._spark_config.items(): spark_builder.config(key, val) diff --git a/basic/pi_spark/tests/test_app.py b/basic/pi_spark/tests/test_app.py index ad57632..6eef572 100644 --- a/basic/pi_spark/tests/test_app.py +++ b/basic/pi_spark/tests/test_app.py @@ -8,5 +8,5 @@ def test_pi_demo_runs(): date_string = "2020-01-01" - result = run(spark, "dev", date_string, 2) + run(spark, "dev", date_string, 2, 10) \ No newline at end of file