{
 "cells": [
  {
   "cell_type": "markdown",
   "metadata": {},
   "source": [
    "# Build Slowly Changing Dimensions Type 2 (SCD2) with Apache Spark and Apache Hudi on Amazon EMR"
   ]
  },
  {
   "cell_type": "markdown",
   "metadata": {},
   "source": [
    "## Prerequisites\n",
    "\n",
    "\n",
    "1. An  [AWS account](https://aws.amazon.com/premiumsupport/knowledge-center/create-and-activate-aws-account/) \n",
    "2. To be able to run the code sample,  [create an EMR cluster](https://docs.aws.amazon.com/emr/latest/ManagementGuide/emr-setting-up.html)  with  [EMR notebook](https://docs.aws.amazon.com/emr/latest/ManagementGuide/emr-managed-notebooks-create.html) . It is important to use the EMR version greate than 5.28.\n",
    "3. To enable Apache Hudi on EMR, run the setup as described in the documentation  [How to use Hudi with Amazon EMR Notebooks](https://docs.aws.amazon.com/emr/latest/ReleaseGuide/emr-hudi-work-with-dataset.html) \n",
    "4. [Create S3 bucket](https://docs.aws.amazon.com/AmazonS3/latest/user-guide/create-bucket.html)  where Hudi files will be stored\n",
    "\n"
   ]
  },
  {
   "cell_type": "markdown",
   "metadata": {},
   "source": [
    "## Enable Hudi libraries"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": 1,
   "metadata": {},
   "outputs": [
    {
     "data": {
      "application/vnd.jupyter.widget-view+json": {
       "model_id": "f70dc69c16b645de9a143137e417c31b",
       "version_major": 2,
       "version_minor": 0
      },
      "text/plain": [
       "VBox()"
      ]
     },
     "metadata": {},
     "output_type": "display_data"
    },
    {
     "name": "stdout",
     "output_type": "stream",
     "text": [
      "Starting Spark application\n"
     ]
    },
    {
     "data": {
      "text/html": [
       "<table>\n",
       "<tr><th>ID</th><th>YARN Application ID</th><th>Kind</th><th>State</th><th>Spark UI</th><th>Driver log</th><th>Current session?</th></tr><tr><td>0</td><td>application_1615560907650_0002</td><td>pyspark</td><td>idle</td><td><a target=\"_blank\" href=\"http://ip-172-31-1-212.eu-west-1.compute.internal:20888/proxy/application_1615560907650_0002/\" >Link</a></td><td><a target=\"_blank\" href=\"http://ip-172-31-4-74.eu-west-1.compute.internal:8042/node/containerlogs/container_1615560907650_0002_01_000001/livy\" >Link</a></td><td>✔</td></tr></table>"
      ],
      "text/plain": [
       "<IPython.core.display.HTML object>"
      ]
     },
     "metadata": {},
     "output_type": "display_data"
    },
    {
     "data": {
      "application/vnd.jupyter.widget-view+json": {
       "model_id": "",
       "version_major": 2,
       "version_minor": 0
      },
      "text/plain": [
       "FloatProgress(value=0.0, bar_style='info', description='Progress:', layout=Layout(height='25px', width='50%'),…"
      ]
     },
     "metadata": {},
     "output_type": "display_data"
    },
    {
     "name": "stdout",
     "output_type": "stream",
     "text": [
      "SparkSession available as 'spark'.\n"
     ]
    },
    {
     "data": {
      "application/vnd.jupyter.widget-view+json": {
       "model_id": "",
       "version_major": 2,
       "version_minor": 0
      },
      "text/plain": [
       "FloatProgress(value=0.0, bar_style='info', description='Progress:', layout=Layout(height='25px', width='50%'),…"
      ]
     },
     "metadata": {},
     "output_type": "display_data"
    },
    {
     "name": "stdout",
     "output_type": "stream",
     "text": [
      "0"
     ]
    }
   ],
   "source": [
    "import subprocess\n",
    "some_path = '/apps/hudi/lib'\n",
    "subprocess.call([\"hdfs\", \"dfs\", \"-mkdir\", \"-p\", some_path])\n",
    "subprocess.call([\"hdfs\", \"dfs\", \"-copyFromLocal\", \"/usr/lib/hudi/hudi-spark-bundle.jar\", \"/apps/hudi/lib/hudi-spark-bundle.jar\"])\n",
    "subprocess.call([\"hdfs\", \"dfs\", \"-chmod\", \"755\", \"/apps/hudi/lib/hudi-spark-bundle.jar\"])"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": 2,
   "metadata": {},
   "outputs": [
    {
     "name": "stdout",
     "output_type": "stream",
     "text": [
      "Starting Spark application\n"
     ]
    },
    {
     "data": {
      "text/html": [
       "<table>\n",
       "<tr><th>ID</th><th>YARN Application ID</th><th>Kind</th><th>State</th><th>Spark UI</th><th>Driver log</th><th>Current session?</th></tr><tr><td>1</td><td>application_1615560907650_0003</td><td>pyspark</td><td>idle</td><td><a target=\"_blank\" href=\"http://ip-172-31-1-212.eu-west-1.compute.internal:20888/proxy/application_1615560907650_0003/\" >Link</a></td><td><a target=\"_blank\" href=\"http://ip-172-31-14-231.eu-west-1.compute.internal:8042/node/containerlogs/container_1615560907650_0003_01_000001/livy\" >Link</a></td><td>✔</td></tr></table>"
      ],
      "text/plain": [
       "<IPython.core.display.HTML object>"
      ]
     },
     "metadata": {},
     "output_type": "display_data"
    },
    {
     "data": {
      "application/vnd.jupyter.widget-view+json": {
       "model_id": "",
       "version_major": 2,
       "version_minor": 0
      },
      "text/plain": [
       "FloatProgress(value=0.0, bar_style='info', description='Progress:', layout=Layout(height='25px', width='50%'),…"
      ]
     },
     "metadata": {},
     "output_type": "display_data"
    },
    {
     "name": "stdout",
     "output_type": "stream",
     "text": [
      "SparkSession available as 'spark'.\n"
     ]
    },
    {
     "data": {
      "text/html": [
       "Current session configs: <tt>{'conf': {'spark.jars': 'hdfs:///apps/hudi/lib/hudi-spark-bundle.jar', 'spark.serializer': 'org.apache.spark.serializer.KryoSerializer', 'spark.sql.hive.convertMetastoreParquet': 'false'}, 'kind': 'pyspark'}</tt><br>"
      ],
      "text/plain": [
       "<IPython.core.display.HTML object>"
      ]
     },
     "metadata": {},
     "output_type": "display_data"
    },
    {
     "data": {
      "text/html": [
       "<table>\n",
       "<tr><th>ID</th><th>YARN Application ID</th><th>Kind</th><th>State</th><th>Spark UI</th><th>Driver log</th><th>Current session?</th></tr><tr><td>1</td><td>application_1615560907650_0003</td><td>pyspark</td><td>idle</td><td><a target=\"_blank\" href=\"http://ip-172-31-1-212.eu-west-1.compute.internal:20888/proxy/application_1615560907650_0003/\" >Link</a></td><td><a target=\"_blank\" href=\"http://ip-172-31-14-231.eu-west-1.compute.internal:8042/node/containerlogs/container_1615560907650_0003_01_000001/livy\" >Link</a></td><td>✔</td></tr></table>"
      ],
      "text/plain": [
       "<IPython.core.display.HTML object>"
      ]
     },
     "metadata": {},
     "output_type": "display_data"
    }
   ],
   "source": [
    "%%configure -f\n",
    "{ \"conf\": {\n",
    "            \"spark.jars\":\"hdfs:///apps/hudi/lib/hudi-spark-bundle.jar\",\n",
    "            \"spark.serializer\":\"org.apache.spark.serializer.KryoSerializer\",\n",
    "            \"spark.sql.hive.convertMetastoreParquet\":\"false\"\n",
    "          }}"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": 3,
   "metadata": {},
   "outputs": [
    {
     "data": {
      "application/vnd.jupyter.widget-view+json": {
       "model_id": "dd0eeff161a14a15b02a98a1634fdbac",
       "version_major": 2,
       "version_minor": 0
      },
      "text/plain": [
       "VBox()"
      ]
     },
     "metadata": {},
     "output_type": "display_data"
    },
    {
     "data": {
      "application/vnd.jupyter.widget-view+json": {
       "model_id": "",
       "version_major": 2,
       "version_minor": 0
      },
      "text/plain": [
       "FloatProgress(value=0.0, bar_style='info', description='Progress:', layout=Layout(height='25px', width='50%'),…"
      ]
     },
     "metadata": {},
     "output_type": "display_data"
    }
   ],
   "source": [
    "from pyspark.sql.types import (\n",
    "    StringType,\n",
    "    StructField,\n",
    "    StructType,\n",
    "    IntegerType,\n",
    "    DoubleType,\n",
    "    DateType,\n",
    "    BooleanType,\n",
    "    TimestampType\n",
    ")"
   ]
  },
  {
   "cell_type": "markdown",
   "metadata": {},
   "source": [
    "## Define customer schema"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": 4,
   "metadata": {},
   "outputs": [
    {
     "data": {
      "application/vnd.jupyter.widget-view+json": {
       "model_id": "20862d98920041acb12c48ddce50251f",
       "version_major": 2,
       "version_minor": 0
      },
      "text/plain": [
       "VBox()"
      ]
     },
     "metadata": {},
     "output_type": "display_data"
    },
    {
     "data": {
      "application/vnd.jupyter.widget-view+json": {
       "model_id": "",
       "version_major": 2,
       "version_minor": 0
      },
      "text/plain": [
       "FloatProgress(value=0.0, bar_style='info', description='Progress:', layout=Layout(height='25px', width='50%'),…"
      ]
     },
     "metadata": {},
     "output_type": "display_data"
    }
   ],
   "source": [
    "dim_customer_schema = StructType([\n",
    "        StructField('customer_id', StringType(), False),\n",
    "        StructField('first_name', StringType(), True),\n",
    "        StructField('last_name', StringType(), True),\n",
    "        StructField('city', StringType(), True),\n",
    "        StructField('country', StringType(), True),\n",
    "        StructField('eff_start_date', DateType(), True),\n",
    "        StructField('eff_end_date', DateType(), True),\n",
    "        StructField('timestamp', TimestampType(), True),\n",
    "        StructField('is_current', BooleanType(), True),\n",
    "    ])"
   ]
  },
  {
   "cell_type": "markdown",
   "metadata": {},
   "source": [
    "## Create customer records"
   ]
  },
  {
   "cell_type": "markdown",
   "metadata": {},
   "source": [
    "For the demo purpose only I will use a current timestamp as a surrogate key for the Customer records"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": 5,
   "metadata": {},
   "outputs": [
    {
     "data": {
      "application/vnd.jupyter.widget-view+json": {
       "model_id": "40fa56a8ce814ac6be96f7e834bc0635",
       "version_major": 2,
       "version_minor": 0
      },
      "text/plain": [
       "VBox()"
      ]
     },
     "metadata": {},
     "output_type": "display_data"
    },
    {
     "data": {
      "application/vnd.jupyter.widget-view+json": {
       "model_id": "",
       "version_major": 2,
       "version_minor": 0
      },
      "text/plain": [
       "FloatProgress(value=0.0, bar_style='info', description='Progress:', layout=Layout(height='25px', width='50%'),…"
      ]
     },
     "metadata": {},
     "output_type": "display_data"
    }
   ],
   "source": [
    "from pyspark.sql.functions import udf\n",
    "import time\n",
    "\n",
    "random_udf = udf(lambda: str(int(time.time() * 1000000)), StringType()) "
   ]
  },
  {
   "cell_type": "code",
   "execution_count": 6,
   "metadata": {},
   "outputs": [
    {
     "data": {
      "application/vnd.jupyter.widget-view+json": {
       "model_id": "5c67fac3df9d4d74957df962583d9729",
       "version_major": 2,
       "version_minor": 0
      },
      "text/plain": [
       "VBox()"
      ]
     },
     "metadata": {},
     "output_type": "display_data"
    },
    {
     "data": {
      "application/vnd.jupyter.widget-view+json": {
       "model_id": "",
       "version_major": 2,
       "version_minor": 0
      },
      "text/plain": [
       "FloatProgress(value=0.0, bar_style='info', description='Progress:', layout=Layout(height='25px', width='50%'),…"
      ]
     },
     "metadata": {},
     "output_type": "display_data"
    },
    {
     "name": "stdout",
     "output_type": "stream",
     "text": [
      "+-----------+----------+---------+-------+-------+--------------+------------+-------------------+----------+----------------+\n",
      "|customer_id|first_name|last_name|city   |country|eff_start_date|eff_end_date|timestamp          |is_current|customer_dim_key|\n",
      "+-----------+----------+---------+-------+-------+--------------+------------+-------------------+----------+----------------+\n",
      "|1          |John      |Smith    |London |UK     |2020-09-27    |2999-12-31  |2020-12-08 09:15:32|true      |1615561177547777|\n",
      "|2          |Susan     |Chas     |Seattle|US     |2020-10-14    |2999-12-31  |2020-12-08 09:15:32|true      |1615561178046248|\n",
      "+-----------+----------+---------+-------+-------+--------------+------------+-------------------+----------+----------------+"
     ]
    }
   ],
   "source": [
    "from datetime import datetime\n",
    "\n",
    "customer_dim_df = spark.createDataFrame([('1', 'John', 'Smith', 'London', 'UK', datetime.strptime('2020-09-27', '%Y-%m-%d'), datetime.strptime('2999-12-31', '%Y-%m-%d'), datetime.strptime('2020-12-08 09:15:32', '%Y-%m-%d %H:%M:%S'), True),\n",
    "                       ('2', 'Susan', 'Chas', 'Seattle', 'US', datetime.strptime('2020-10-14', '%Y-%m-%d'), datetime.strptime('2999-12-31', '%Y-%m-%d'), datetime.strptime('2020-12-08 09:15:32', '%Y-%m-%d %H:%M:%S'), True)], dim_customer_schema)\n",
    "\n",
    "customer_hudi_df = customer_dim_df.withColumn(\"customer_dim_key\", random_udf())\n",
    "\n",
    "customer_hudi_df.cache()\n",
    "\n",
    "customer_hudi_df.show(5, False)"
   ]
  },
  {
   "cell_type": "markdown",
   "metadata": {},
   "source": [
    "## Store customer records"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": 7,
   "metadata": {},
   "outputs": [
    {
     "data": {
      "application/vnd.jupyter.widget-view+json": {
       "model_id": "186a0cae022d48779d17ebe79593f224",
       "version_major": 2,
       "version_minor": 0
      },
      "text/plain": [
       "VBox()"
      ]
     },
     "metadata": {},
     "output_type": "display_data"
    },
    {
     "data": {
      "application/vnd.jupyter.widget-view+json": {
       "model_id": "",
       "version_major": 2,
       "version_minor": 0
      },
      "text/plain": [
       "FloatProgress(value=0.0, bar_style='info', description='Progress:', layout=Layout(height='25px', width='50%'),…"
      ]
     },
     "metadata": {},
     "output_type": "display_data"
    }
   ],
   "source": [
    "customers_table_path = 's3://MY-BUCKET/customers/'\n",
    "customers_table_name = 'customers_table'\n",
    "\n",
    "partition_key = \"country:SIMPLE\"\n",
    "\n",
    "hudi_options = {'hoodie.insert.shuffle.parallelism':'2',\n",
    "               'hoodie.upsert.shuffle.parallelism':'2',\n",
    "               'hoodie.delete.shuffle.parallelism':'2',\n",
    "               'hoodie.bulkinsert.shuffle.parallelism':'2',\n",
    "               'hoodie.datasource.hive_sync.enable':'false',\n",
    "               'hoodie.datasource.hive_sync.assume_date_partitioning':'true'\n",
    "}\n",
    "\n",
    "customer_hudi_df.write.format('org.apache.hudi')\\\n",
    "                             .options(**hudi_options)\\\n",
    "                             .option('hoodie.table.name',customers_table_name)\\\n",
    "                             .option('hoodie.datasource.write.recordkey.field','customer_dim_key')\\\n",
    "                             .option('hoodie.datasource.write.partitionpath.field',partition_key)\\\n",
    "                             .option('hoodie.datasource.write.precombine.field','timestamp')\\\n",
    "                             .option('hoodie.datasource.write.operation', 'insert')\\\n",
    "                             .option('hoodie.datasource.hive_sync.table',customers_table_name)\\\n",
    "                             .mode('Overwrite')\\\n",
    "                             .save(customers_table_path)"
   ]
  },
  {
   "cell_type": "markdown",
   "metadata": {},
   "source": [
    "## Read customers Hudi table"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": 8,
   "metadata": {},
   "outputs": [
    {
     "data": {
      "application/vnd.jupyter.widget-view+json": {
       "model_id": "f2a142c69527476e96585d8c327beecc",
       "version_major": 2,
       "version_minor": 0
      },
      "text/plain": [
       "VBox()"
      ]
     },
     "metadata": {},
     "output_type": "display_data"
    },
    {
     "data": {
      "application/vnd.jupyter.widget-view+json": {
       "model_id": "",
       "version_major": 2,
       "version_minor": 0
      },
      "text/plain": [
       "FloatProgress(value=0.0, bar_style='info', description='Progress:', layout=Layout(height='25px', width='50%'),…"
      ]
     },
     "metadata": {},
     "output_type": "display_data"
    },
    {
     "name": "stdout",
     "output_type": "stream",
     "text": [
      "+-----------+----------+---------+-------+-------+--------------+------------+----------------+----------+\n",
      "|customer_id|first_name|last_name|city   |country|eff_start_date|eff_end_date|customer_dim_key|is_current|\n",
      "+-----------+----------+---------+-------+-------+--------------+------------+----------------+----------+\n",
      "|1          |John      |Smith    |London |UK     |2020-09-27    |2999-12-31  |1615561177547777|true      |\n",
      "|2          |Susan     |Chas     |Seattle|US     |2020-10-14    |2999-12-31  |1615561178046248|true      |\n",
      "+-----------+----------+---------+-------+-------+--------------+------------+----------------+----------+"
     ]
    }
   ],
   "source": [
    "customer_hudi_df = spark. \\\n",
    "  read. \\\n",
    "  format(\"hudi\"). \\\n",
    "  load(customers_table_path + 'default/')\n",
    "\n",
    "customer_hudi_df.createOrReplaceTempView(customers_table_name)\n",
    "\n",
    "spark.sql('select customer_id,'\n",
    "          'first_name, '\n",
    "          'last_name, '\n",
    "          'city, '\n",
    "          'country, '\n",
    "          'eff_start_date, '\n",
    "          'eff_end_date, '\n",
    "          'customer_dim_key, '\n",
    "          'is_current '\n",
    "          'from ' + customers_table_name).show(3, False)"
   ]
  },
  {
   "cell_type": "markdown",
   "metadata": {},
   "source": [
    "## Create sales"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": 9,
   "metadata": {},
   "outputs": [
    {
     "data": {
      "application/vnd.jupyter.widget-view+json": {
       "model_id": "5aa8fa6b59b34b0fb45fea4c39eb0360",
       "version_major": 2,
       "version_minor": 0
      },
      "text/plain": [
       "VBox()"
      ]
     },
     "metadata": {},
     "output_type": "display_data"
    },
    {
     "data": {
      "application/vnd.jupyter.widget-view+json": {
       "model_id": "",
       "version_major": 2,
       "version_minor": 0
      },
      "text/plain": [
       "FloatProgress(value=0.0, bar_style='info', description='Progress:', layout=Layout(height='25px', width='50%'),…"
      ]
     },
     "metadata": {},
     "output_type": "display_data"
    },
    {
     "name": "stdout",
     "output_type": "stream",
     "text": [
      "+-------+--------+------+-------------------+-----------+\n",
      "|item_id|quantity| price|          timestamp|customer_id|\n",
      "+-------+--------+------+-------------------+-----------+\n",
      "|    100|      25|123.46|2020-11-17 09:15:32|          1|\n",
      "|    101|     300|123.46|2020-10-28 09:15:32|          1|\n",
      "|    102|       5|1038.0|2020-12-08 09:15:32|          2|\n",
      "+-------+--------+------+-------------------+-----------+"
     ]
    }
   ],
   "source": [
    "from pyspark.sql.functions import to_timestamp\n",
    "\n",
    "fact_sales_schema = StructType([\n",
    "        StructField('item_id', StringType(), True),\n",
    "        StructField('quantity', IntegerType(), True),\n",
    "        StructField('price', DoubleType(), True),\n",
    "        StructField('timestamp', TimestampType(), True),\n",
    "        StructField('customer_id', StringType(), True)\n",
    "    ])\n",
    "\n",
    "sales_fact_df = spark.createDataFrame([('100', 25, 123.46, datetime.strptime('2020-11-17 09:15:32', '%Y-%m-%d %H:%M:%S'), '1'),\n",
    "                                       ('101', 300, 123.46, datetime.strptime('2020-10-28 09:15:32', '%Y-%m-%d %H:%M:%S'), '1'),\n",
    "                                      ('102', 5, 1038.0, datetime.strptime('2020-12-08 09:15:32', '%Y-%m-%d %H:%M:%S'), '2')], fact_sales_schema)\n",
    "\n",
    "sales_fact_df.show()"
   ]
  },
  {
   "cell_type": "markdown",
   "metadata": {},
   "source": [
    "## Customer dimension key lookup\n",
    "\n"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": 10,
   "metadata": {},
   "outputs": [
    {
     "data": {
      "application/vnd.jupyter.widget-view+json": {
       "model_id": "67f126beb7c7469fa7ad04c1ab41096a",
       "version_major": 2,
       "version_minor": 0
      },
      "text/plain": [
       "VBox()"
      ]
     },
     "metadata": {},
     "output_type": "display_data"
    },
    {
     "data": {
      "application/vnd.jupyter.widget-view+json": {
       "model_id": "",
       "version_major": 2,
       "version_minor": 0
      },
      "text/plain": [
       "FloatProgress(value=0.0, bar_style='info', description='Progress:', layout=Layout(height='25px', width='50%'),…"
      ]
     },
     "metadata": {},
     "output_type": "display_data"
    }
   ],
   "source": [
    "from pyspark.sql.functions import when\n",
    "\n",
    "join_cond = [sales_fact_df.customer_id == customer_hudi_df.customer_id,\n",
    "             sales_fact_df.timestamp >= customer_hudi_df.eff_start_date,\n",
    "             sales_fact_df.timestamp < customer_hudi_df.eff_end_date]\n",
    "\n",
    "customers_dim_key_df = (sales_fact_df\n",
    "                          .join(customer_hudi_df, join_cond, 'leftouter')\n",
    "                          .select(sales_fact_df['*'],\n",
    "                            when(customer_hudi_df.customer_dim_key.isNull(), '-1')\n",
    "                                  .otherwise(customer_hudi_df.customer_dim_key)\n",
    "                                  .alias(\"customer_dim_key\") )\n",
    "                       )"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": 11,
   "metadata": {},
   "outputs": [
    {
     "data": {
      "application/vnd.jupyter.widget-view+json": {
       "model_id": "b5c915194756499d8d5710ca1f27686d",
       "version_major": 2,
       "version_minor": 0
      },
      "text/plain": [
       "VBox()"
      ]
     },
     "metadata": {},
     "output_type": "display_data"
    },
    {
     "data": {
      "application/vnd.jupyter.widget-view+json": {
       "model_id": "",
       "version_major": 2,
       "version_minor": 0
      },
      "text/plain": [
       "FloatProgress(value=0.0, bar_style='info', description='Progress:', layout=Layout(height='25px', width='50%'),…"
      ]
     },
     "metadata": {},
     "output_type": "display_data"
    },
    {
     "name": "stdout",
     "output_type": "stream",
     "text": [
      "+-------+--------+------+-------------------+-----------+----------------+\n",
      "|item_id|quantity| price|          timestamp|customer_id|customer_dim_key|\n",
      "+-------+--------+------+-------------------+-----------+----------------+\n",
      "|    100|      25|123.46|2020-11-17 09:15:32|          1|1615561177547777|\n",
      "|    101|     300|123.46|2020-10-28 09:15:32|          1|1615561177547777|\n",
      "|    102|       5|1038.0|2020-12-08 09:15:32|          2|1615561178046248|\n",
      "+-------+--------+------+-------------------+-----------+----------------+"
     ]
    }
   ],
   "source": [
    "customers_dim_key_df.show()"
   ]
  },
  {
   "cell_type": "markdown",
   "metadata": {},
   "source": [
    "## Save sales"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": 12,
   "metadata": {},
   "outputs": [
    {
     "data": {
      "application/vnd.jupyter.widget-view+json": {
       "model_id": "702a92a23781483b89d78f15ec97d694",
       "version_major": 2,
       "version_minor": 0
      },
      "text/plain": [
       "VBox()"
      ]
     },
     "metadata": {},
     "output_type": "display_data"
    },
    {
     "data": {
      "application/vnd.jupyter.widget-view+json": {
       "model_id": "",
       "version_major": 2,
       "version_minor": 0
      },
      "text/plain": [
       "FloatProgress(value=0.0, bar_style='info', description='Progress:', layout=Layout(height='25px', width='50%'),…"
      ]
     },
     "metadata": {},
     "output_type": "display_data"
    }
   ],
   "source": [
    "sales_table_path = 's3://MY-BUCKET/sales/'\n",
    "sales_table_name = 'sales_table'\n",
    "\n",
    "customers_dim_key_df.write.format('parquet')\\\n",
    "                             .mode('Overwrite')\\\n",
    "                             .save(sales_table_path)"
   ]
  },
  {
   "cell_type": "markdown",
   "metadata": {},
   "source": [
    "## Read sales table"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": 13,
   "metadata": {},
   "outputs": [
    {
     "data": {
      "application/vnd.jupyter.widget-view+json": {
       "model_id": "26d8402ae2a447d398a9bb76480911d1",
       "version_major": 2,
       "version_minor": 0
      },
      "text/plain": [
       "VBox()"
      ]
     },
     "metadata": {},
     "output_type": "display_data"
    },
    {
     "data": {
      "application/vnd.jupyter.widget-view+json": {
       "model_id": "",
       "version_major": 2,
       "version_minor": 0
      },
      "text/plain": [
       "FloatProgress(value=0.0, bar_style='info', description='Progress:', layout=Layout(height='25px', width='50%'),…"
      ]
     },
     "metadata": {},
     "output_type": "display_data"
    },
    {
     "name": "stdout",
     "output_type": "stream",
     "text": [
      "+-------+--------+------+-------------------+-----------+----------------+\n",
      "|item_id|quantity|price |timestamp          |customer_id|customer_dim_key|\n",
      "+-------+--------+------+-------------------+-----------+----------------+\n",
      "|100    |25      |123.46|2020-11-17 09:15:32|1          |1615561177547777|\n",
      "|101    |300     |123.46|2020-10-28 09:15:32|1          |1615561177547777|\n",
      "|102    |5       |1038.0|2020-12-08 09:15:32|2          |1615561178046248|\n",
      "+-------+--------+------+-------------------+-----------+----------------+"
     ]
    }
   ],
   "source": [
    "sales_hudi_df = spark. \\\n",
    "  read. \\\n",
    "  load(sales_table_path)\n",
    "\n",
    "sales_hudi_df.createOrReplaceTempView(sales_table_name)\n",
    "\n",
    "spark.sql('select item_id, '\n",
    "          'quantity,'\n",
    "          'price,'\n",
    "          'timestamp,'\n",
    "          'customer_id,'\n",
    "          'customer_dim_key from ' + sales_table_name).show(5, False)"
   ]
  },
  {
   "cell_type": "markdown",
   "metadata": {},
   "source": [
    "## Get number of sales per country"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": 14,
   "metadata": {},
   "outputs": [
    {
     "data": {
      "application/vnd.jupyter.widget-view+json": {
       "model_id": "f20e0a9c25b84ffeb27dcf407b6653eb",
       "version_major": 2,
       "version_minor": 0
      },
      "text/plain": [
       "VBox()"
      ]
     },
     "metadata": {},
     "output_type": "display_data"
    },
    {
     "data": {
      "application/vnd.jupyter.widget-view+json": {
       "model_id": "",
       "version_major": 2,
       "version_minor": 0
      },
      "text/plain": [
       "FloatProgress(value=0.0, bar_style='info', description='Progress:', layout=Layout(height='25px', width='50%'),…"
      ]
     },
     "metadata": {},
     "output_type": "display_data"
    },
    {
     "name": "stdout",
     "output_type": "stream",
     "text": [
      "+-------+--------------+-----------+\n",
      "|country|sales_quantity|count_sales|\n",
      "+-------+--------------+-----------+\n",
      "|     US|             5|          1|\n",
      "|     UK|           325|          2|\n",
      "+-------+--------------+-----------+"
     ]
    }
   ],
   "source": [
    "spark.sql(\n",
    "    'SELECT ct.country, '\n",
    "    'SUM(st.quantity) as sales_quantity,'\n",
    "    'COUNT(*) as count_sales '\n",
    "    'FROM sales_table st '\n",
    "    'INNER JOIN customers_table ct on st.customer_dim_key = ct.customer_dim_key '\n",
    "    'group by ct.country').show()"
   ]
  },
  {
   "cell_type": "markdown",
   "metadata": {},
   "source": [
    "## Customer Susan changed the country from US to FR, new customer Bastian created"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": 15,
   "metadata": {},
   "outputs": [
    {
     "data": {
      "application/vnd.jupyter.widget-view+json": {
       "model_id": "b1602adc1d4545ddba494e4af7c3570c",
       "version_major": 2,
       "version_minor": 0
      },
      "text/plain": [
       "VBox()"
      ]
     },
     "metadata": {},
     "output_type": "display_data"
    },
    {
     "data": {
      "application/vnd.jupyter.widget-view+json": {
       "model_id": "",
       "version_major": 2,
       "version_minor": 0
      },
      "text/plain": [
       "FloatProgress(value=0.0, bar_style='info', description='Progress:', layout=Layout(height='25px', width='50%'),…"
      ]
     },
     "metadata": {},
     "output_type": "display_data"
    },
    {
     "name": "stdout",
     "output_type": "stream",
     "text": [
      "+-----------+----------+---------+------+-------+--------------+------------+-------------------+----------+----------------+\n",
      "|customer_id|first_name|last_name|  city|country|eff_start_date|eff_end_date|          timestamp|is_current|customer_dim_key|\n",
      "+-----------+----------+---------+------+-------+--------------+------------+-------------------+----------+----------------+\n",
      "|          3|   Bastian|     Back|Berlin|     GE|    2021-03-12|  2999-12-31|2020-12-09 09:15:32|      true|1615561211194211|\n",
      "|          2|     Susan|     Chas| Paris|     FR|    2021-03-12|  2999-12-31|2020-12-09 10:15:32|      true|1615561211350200|\n",
      "+-----------+----------+---------+------+-------+--------------+------------+-------------------+----------+----------------+"
     ]
    }
   ],
   "source": [
    "new_customer_dim_df = spark.createDataFrame([('3', 'Bastian', 'Back', \n",
    "                    'Berlin', 'GE',\n",
    "                    datetime.strptime(datetime.today().strftime('%Y-%m-%d'), '%Y-%m-%d'),\n",
    "                    datetime.strptime('2999-12-31', '%Y-%m-%d'), \n",
    "                    datetime.strptime('2020-12-09 09:15:32', '%Y-%m-%d %H:%M:%S'), True),\n",
    "                    ('2', 'Susan', 'Chas',\n",
    "                    'Paris', 'FR',\n",
    "                    datetime.strptime(datetime.today().strftime('%Y-%m-%d'), '%Y-%m-%d'),\n",
    "                    datetime.strptime('2999-12-31', '%Y-%m-%d'), \n",
    "                    datetime.strptime('2020-12-09 10:15:32', '%Y-%m-%d %H:%M:%S'), True)],\n",
    "                dim_customer_schema)\n",
    "\n",
    "new_customer_dim_df = new_customer_dim_df.withColumn(\"customer_dim_key\", random_udf())\n",
    "\n",
    "new_customer_dim_df.cache()\n",
    "\n",
    "new_customer_dim_df.show()"
   ]
  },
  {
   "cell_type": "markdown",
   "metadata": {},
   "source": [
    "## Customers UPSERT"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": 16,
   "metadata": {},
   "outputs": [
    {
     "data": {
      "application/vnd.jupyter.widget-view+json": {
       "model_id": "7845f48ad9f0461c8d0ff5bbc83aee99",
       "version_major": 2,
       "version_minor": 0
      },
      "text/plain": [
       "VBox()"
      ]
     },
     "metadata": {},
     "output_type": "display_data"
    },
    {
     "data": {
      "application/vnd.jupyter.widget-view+json": {
       "model_id": "",
       "version_major": 2,
       "version_minor": 0
      },
      "text/plain": [
       "FloatProgress(value=0.0, bar_style='info', description='Progress:', layout=Layout(height='25px', width='50%'),…"
      ]
     },
     "metadata": {},
     "output_type": "display_data"
    }
   ],
   "source": [
    "from pyspark.sql.functions import lit\n",
    "\n",
    "join_cond = [customer_hudi_df.customer_id == new_customer_dim_df.customer_id, customer_hudi_df.is_current == True]\n",
    "\n",
    "## Find customer records to update\n",
    "customers_to_update_df = (customer_hudi_df\n",
    "                          .join(new_customer_dim_df, join_cond)\n",
    "                          .select(customer_hudi_df.customer_id,\n",
    "                                  customer_hudi_df.first_name,\n",
    "                                  customer_hudi_df.last_name,\n",
    "                                  customer_hudi_df.city,\n",
    "                                  customer_hudi_df.country,\n",
    "                                  customer_hudi_df.eff_start_date,\n",
    "                                  new_customer_dim_df.eff_start_date.alias(\"eff_end_date\"),\n",
    "                                  customer_hudi_df.customer_dim_key,\n",
    "                                  customer_hudi_df.timestamp)\n",
    "                          .withColumn('is_current', lit(False))\n",
    "                         )\n",
    "\n",
    "\n",
    "## Union with new customer records\n",
    "merged_customers_df = new_customer_dim_df.unionByName(customers_to_update_df)\n",
    "\n",
    "partition_key = \"country:SIMPLE\"\n",
    "\n",
    "# Upsert\n",
    "merged_customers_df.write.format('org.apache.hudi')\\\n",
    "                    .options(**hudi_options)\\\n",
    "                    .option('hoodie.table.name',customers_table_name)\\\n",
    "                    .option('hoodie.datasource.write.recordkey.field','customer_dim_key')\\\n",
    "                    .option('hoodie.datasource.write.partitionpath.field',partition_key)\\\n",
    "                    .option('hoodie.datasource.write.precombine.field','timestamp')\\\n",
    "                    .option('hoodie.datasource.write.operation', 'upsert')\\\n",
    "                    .option('hoodie.datasource.hive_sync.table',customers_table_name)\\\n",
    "                    .mode('append')\\\n",
    "                    .save(customers_table_path)"
   ]
  },
  {
   "cell_type": "markdown",
   "metadata": {},
   "source": [
    "## Read customers Hudi table"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": 17,
   "metadata": {},
   "outputs": [
    {
     "data": {
      "application/vnd.jupyter.widget-view+json": {
       "model_id": "60bd43f760cd433486a2e4372d3287d2",
       "version_major": 2,
       "version_minor": 0
      },
      "text/plain": [
       "VBox()"
      ]
     },
     "metadata": {},
     "output_type": "display_data"
    },
    {
     "data": {
      "application/vnd.jupyter.widget-view+json": {
       "model_id": "",
       "version_major": 2,
       "version_minor": 0
      },
      "text/plain": [
       "FloatProgress(value=0.0, bar_style='info', description='Progress:', layout=Layout(height='25px', width='50%'),…"
      ]
     },
     "metadata": {},
     "output_type": "display_data"
    },
    {
     "name": "stdout",
     "output_type": "stream",
     "text": [
      "+-----------+----------+---------+-------+-------+--------------+------------+----------------+----------+\n",
      "|customer_id|first_name|last_name|   city|country|eff_start_date|eff_end_date|customer_dim_key|is_current|\n",
      "+-----------+----------+---------+-------+-------+--------------+------------+----------------+----------+\n",
      "|          1|      John|    Smith| London|     UK|    2020-09-27|  2999-12-31|1615561177547777|      true|\n",
      "|          2|     Susan|     Chas|Seattle|     US|    2020-10-14|  2021-03-12|1615561178046248|     false|\n",
      "|          2|     Susan|     Chas|  Paris|     FR|    2021-03-12|  2999-12-31|1615561211350200|      true|\n",
      "|          3|   Bastian|     Back| Berlin|     GE|    2021-03-12|  2999-12-31|1615561211194211|      true|\n",
      "+-----------+----------+---------+-------+-------+--------------+------------+----------------+----------+"
     ]
    }
   ],
   "source": [
    "customer_hudi_df = spark. \\\n",
    "  read. \\\n",
    "  format(\"hudi\"). \\\n",
    "  load(customers_table_path + 'default/')\n",
    "\n",
    "customer_hudi_df.createOrReplaceTempView(customers_table_name)\n",
    "\n",
    "spark.sql('select customer_id,'\n",
    "          'first_name, '\n",
    "          'last_name, '\n",
    "          'city, '\n",
    "          'country, '\n",
    "          'eff_start_date, '\n",
    "          'eff_end_date, '\n",
    "          'customer_dim_key, '\n",
    "          'is_current '\n",
    "          'from ' + customers_table_name).show()"
   ]
  },
  {
   "cell_type": "markdown",
   "metadata": {},
   "source": [
    "## Add new sales for Susan"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": 18,
   "metadata": {},
   "outputs": [
    {
     "data": {
      "application/vnd.jupyter.widget-view+json": {
       "model_id": "a939c0a0d7764daeb0958a265b86a0d3",
       "version_major": 2,
       "version_minor": 0
      },
      "text/plain": [
       "VBox()"
      ]
     },
     "metadata": {},
     "output_type": "display_data"
    },
    {
     "data": {
      "application/vnd.jupyter.widget-view+json": {
       "model_id": "",
       "version_major": 2,
       "version_minor": 0
      },
      "text/plain": [
       "FloatProgress(value=0.0, bar_style='info', description='Progress:', layout=Layout(height='25px', width='50%'),…"
      ]
     },
     "metadata": {},
     "output_type": "display_data"
    },
    {
     "name": "stdout",
     "output_type": "stream",
     "text": [
      "+-------+--------+------+-------------------+-----------+\n",
      "|item_id|quantity| price|          timestamp|customer_id|\n",
      "+-------+--------+------+-------------------+-----------+\n",
      "|    103|     250|  12.3|2021-03-12 12:15:42|          2|\n",
      "|    104|       3|1021.0|2021-03-12 06:35:32|          2|\n",
      "+-------+--------+------+-------------------+-----------+"
     ]
    }
   ],
   "source": [
    "sales_fact_df = spark.createDataFrame([('103', 250, 12.3, datetime.strptime(datetime.today().strftime('%Y-%m-%d')+' 12:15:42', '%Y-%m-%d %H:%M:%S'), '2'),\n",
    "                                       ('104', 3, 1021.0, datetime.strptime(datetime.today().strftime('%Y-%m-%d')+' 06:35:32', '%Y-%m-%d %H:%M:%S'), '2')], fact_sales_schema)\n",
    "\n",
    "sales_fact_df.show()"
   ]
  },
  {
   "cell_type": "markdown",
   "metadata": {},
   "source": [
    "## Customer dimention key lookup\n"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": 19,
   "metadata": {},
   "outputs": [
    {
     "data": {
      "application/vnd.jupyter.widget-view+json": {
       "model_id": "772bde4f8e07452da6e21f027fc90781",
       "version_major": 2,
       "version_minor": 0
      },
      "text/plain": [
       "VBox()"
      ]
     },
     "metadata": {},
     "output_type": "display_data"
    },
    {
     "data": {
      "application/vnd.jupyter.widget-view+json": {
       "model_id": "",
       "version_major": 2,
       "version_minor": 0
      },
      "text/plain": [
       "FloatProgress(value=0.0, bar_style='info', description='Progress:', layout=Layout(height='25px', width='50%'),…"
      ]
     },
     "metadata": {},
     "output_type": "display_data"
    },
    {
     "name": "stdout",
     "output_type": "stream",
     "text": [
      "+-------+--------+------+-------------------+-----------+----------------+\n",
      "|item_id|quantity| price|          timestamp|customer_id|customer_dim_key|\n",
      "+-------+--------+------+-------------------+-----------+----------------+\n",
      "|    103|     250|  12.3|2021-03-12 12:15:42|          2|1615561211350200|\n",
      "|    104|       3|1021.0|2021-03-12 06:35:32|          2|1615561211350200|\n",
      "+-------+--------+------+-------------------+-----------+----------------+"
     ]
    }
   ],
   "source": [
    "from pyspark.sql.functions import when\n",
    "\n",
    "join_cond = [sales_fact_df.customer_id == customer_hudi_df.customer_id, sales_fact_df.timestamp >= customer_hudi_df.eff_start_date, sales_fact_df.timestamp < customer_hudi_df.eff_end_date]\n",
    "\n",
    "\n",
    "customers_dim_key_df = (sales_fact_df\n",
    "                          .join(customer_hudi_df, join_cond, 'leftouter')\n",
    "                          .select(sales_fact_df['*'],\n",
    "                            when(customer_hudi_df.customer_dim_key.isNull(), '-1').otherwise(customer_hudi_df.customer_dim_key).alias(\"customer_dim_key\") )\n",
    "                         )\n",
    "\n",
    "customers_dim_key_df.show()"
   ]
  },
  {
   "cell_type": "markdown",
   "metadata": {},
   "source": [
    "## Store new sales"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": 20,
   "metadata": {},
   "outputs": [
    {
     "data": {
      "application/vnd.jupyter.widget-view+json": {
       "model_id": "eb3add58bca94d27b451bba7bdd7159e",
       "version_major": 2,
       "version_minor": 0
      },
      "text/plain": [
       "VBox()"
      ]
     },
     "metadata": {},
     "output_type": "display_data"
    },
    {
     "data": {
      "application/vnd.jupyter.widget-view+json": {
       "model_id": "",
       "version_major": 2,
       "version_minor": 0
      },
      "text/plain": [
       "FloatProgress(value=0.0, bar_style='info', description='Progress:', layout=Layout(height='25px', width='50%'),…"
      ]
     },
     "metadata": {},
     "output_type": "display_data"
    }
   ],
   "source": [
    "sales_table_path = 's3://MY-BUCKET/sales/'\n",
    "sales_table_name = 'sales_table'\n",
    "\n",
    "customers_dim_key_df.write.format('parquet')\\\n",
    "                             .mode('Append')\\\n",
    "                             .save(sales_table_path)"
   ]
  },
  {
   "cell_type": "markdown",
   "metadata": {},
   "source": [
    "## Read sales table"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": 21,
   "metadata": {},
   "outputs": [
    {
     "data": {
      "application/vnd.jupyter.widget-view+json": {
       "model_id": "5f0fc1bba23f489bb44028a1c1704db7",
       "version_major": 2,
       "version_minor": 0
      },
      "text/plain": [
       "VBox()"
      ]
     },
     "metadata": {},
     "output_type": "display_data"
    },
    {
     "data": {
      "application/vnd.jupyter.widget-view+json": {
       "model_id": "",
       "version_major": 2,
       "version_minor": 0
      },
      "text/plain": [
       "FloatProgress(value=0.0, bar_style='info', description='Progress:', layout=Layout(height='25px', width='50%'),…"
      ]
     },
     "metadata": {},
     "output_type": "display_data"
    },
    {
     "name": "stdout",
     "output_type": "stream",
     "text": [
      "+-------+--------+------+-------------------+-----------+----------------+\n",
      "|item_id|quantity|price |timestamp          |customer_id|customer_dim_key|\n",
      "+-------+--------+------+-------------------+-----------+----------------+\n",
      "|100    |25      |123.46|2020-11-17 09:15:32|1          |1615561177547777|\n",
      "|103    |250     |12.3  |2021-03-12 12:15:42|2          |1615561211350200|\n",
      "|101    |300     |123.46|2020-10-28 09:15:32|1          |1615561177547777|\n",
      "|102    |5       |1038.0|2020-12-08 09:15:32|2          |1615561178046248|\n",
      "|104    |3       |1021.0|2021-03-12 06:35:32|2          |1615561211350200|\n",
      "+-------+--------+------+-------------------+-----------+----------------+"
     ]
    }
   ],
   "source": [
    "sales_hudi_df = spark. \\\n",
    "  read. \\\n",
    "  load(sales_table_path)\n",
    "\n",
    "sales_hudi_df.createOrReplaceTempView(sales_table_name)\n",
    "\n",
    "spark.sql(\"select item_id, quantity, price, timestamp, customer_id, customer_dim_key from \" + sales_table_name).show(5, False)"
   ]
  },
  {
   "cell_type": "markdown",
   "metadata": {},
   "source": [
    "## Get number of sales per country"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": 22,
   "metadata": {},
   "outputs": [
    {
     "data": {
      "application/vnd.jupyter.widget-view+json": {
       "model_id": "3cc35d39360d4fc88c315e66676798fa",
       "version_major": 2,
       "version_minor": 0
      },
      "text/plain": [
       "VBox()"
      ]
     },
     "metadata": {},
     "output_type": "display_data"
    },
    {
     "data": {
      "application/vnd.jupyter.widget-view+json": {
       "model_id": "",
       "version_major": 2,
       "version_minor": 0
      },
      "text/plain": [
       "FloatProgress(value=0.0, bar_style='info', description='Progress:', layout=Layout(height='25px', width='50%'),…"
      ]
     },
     "metadata": {},
     "output_type": "display_data"
    },
    {
     "name": "stdout",
     "output_type": "stream",
     "text": [
      "+-------+--------------+-----------+\n",
      "|country|sales_quantity|count_sales|\n",
      "+-------+--------------+-----------+\n",
      "|     US|             5|          1|\n",
      "|     UK|           325|          2|\n",
      "|     FR|           253|          2|\n",
      "+-------+--------------+-----------+"
     ]
    }
   ],
   "source": [
    "spark.sql(\n",
    "    'SELECT ct.country, SUM(st.quantity) as sales_quantity, COUNT(*) as count_sales '\n",
    "    'FROM sales_table st '\n",
    "    'INNER JOIN customers_table ct on st.customer_dim_key = ct.customer_dim_key group by ct.country').show()"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": null,
   "metadata": {},
   "outputs": [],
   "source": []
  }
 ],
 "metadata": {
  "kernelspec": {
   "display_name": "PySpark",
   "language": "",
   "name": "pysparkkernel"
  },
  "language_info": {
   "codemirror_mode": {
    "name": "python",
    "version": 3
   },
   "mimetype": "text/x-python",
   "name": "pyspark",
   "pygments_lexer": "python3"
  }
 },
 "nbformat": 4,
 "nbformat_minor": 4
}
