Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 2 additions & 1 deletion .gitignore
Original file line number Diff line number Diff line change
@@ -1,2 +1,3 @@
.idea/
.conveyor
.conveyor
.terraform/
2 changes: 1 addition & 1 deletion basic/pi_spark/Dockerfile
Original file line number Diff line number Diff line change
@@ -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
Expand Down
34 changes: 27 additions & 7 deletions basic/pi_spark/dags/pi_spark.py
Original file line number Diff line number Diff line change
@@ -1,40 +1,60 @@
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,
application="local:///opt/spark/work-dir/src/pi_spark/app.py",
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",
],
)
2 changes: 1 addition & 1 deletion basic/pi_spark/dev-requirements.in
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
-c requirements.txt
pyspark==3.5.1
pyspark==3.5.6
pandas==1.5.2
pytest>=6.0
pytest-cov
Expand Down
15 changes: 11 additions & 4 deletions basic/pi_spark/dev-requirements.txt
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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
Expand All @@ -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
18 changes: 0 additions & 18 deletions basic/pi_spark/ide.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -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
56 changes: 39 additions & 17 deletions basic/pi_spark/src/pi_spark/app.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,47 +10,69 @@

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
"""
# 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()
main()
1 change: 0 additions & 1 deletion basic/pi_spark/src/pi_spark/common/spark.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
2 changes: 1 addition & 1 deletion basic/pi_spark/tests/test_app.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)