How to manage mangled data when importing from your source in sqoop or pyspark

Viewed 96

I have been working on a project to import the Danish 2.5Million ATM transaction data set to derive some visualizations.

The data is hosted on a mysql server provided by the university. The objective is to import the data using Sqoop and then apply a few transformations to it using pyspark.

Link to the dataset here : https://www.kaggle.com/sparnord/danish-atm-transactions

The Sql server, that hosts this information has a few rows which are intentionally or unintentionally mangled.

Example: So I have a very basic sqoop command which gets the details from the source database. However I run into an issue where there are values which have a double-quote " especially in the column message_text

Sqoop Command :

sqoop import --connect jdbc:mysql:{source-connection-string} --table SRC_ATM_TRANS --username {username}--password {password} --target-dir /user/root/etl_project --fields-terminated-by '|' --lines-terminated-by "\n" -m 1

Here is sample row that is imported in the transaction.

2017|January|1|Sunday|21|Active|85|Diebold Nixdorf|København|Regnbuepladsen|5|1550|55.676|12.571|DKK|MasterCard|4531|Withdrawal|4017|"Suspected malfunction|0.000|55.676|13|2618425|0.000|277|1010|93|3|280.000|0|75|803|Clouds

However the expected output should be

2017|January|1|Sunday|21|Active|85|Diebold Nixdorf|København|Regnbuepladsen|5|1550|55.676|12.571|DKK|MasterCard|4531|Withdrawal|4017|"Suspected malfunction,0.000|55.676|13|2618425|0.000|277|1010|93|3|280.000|0|75|803|Clouds|Cloudy

At first I was okay with this hoping that pyspark would handle the mangled data since the delimiters are specified.

But now I run into issues when populating my dataframe.

transactions = spark.read.option("sep","|").csv("/user/root/etl_project/part-m-00000", header = False,schema = transaction_schema)

However when I inspect my rows I see that the mangled data has caused the dataframe to put these affected values into a single column!

transactions.filter(transactions.message_code == "4017").collect()

Row(year=2017, month=u'January', day=1, weekday=u'Sunday', hour=17, atm_status=u'Active', atm_id=u'35', atm_manufacturer=u'NCR', atm_location=u'Aabybro', atm_streetname=u'\xc3\u0192\xcb\u0153stergade', atm_street_number=6, atm_zipcode=9440, atm_lat=57.162, atm_lon=9.73, currency=u'DKK', card_type=u'MasterCard', transaction_amount=7387, service=u'Withdrawal', message_code=u'4017', message_text=u'Suspected malfunction|0.000|57.158|10|2625037|0.000|276|1021|83|4|319.000|0|0|800|Clear', weather_lat=None, weather_lon=None, weather_city_id=None, weather_city_name=None, temp=None, pressure=None, humidity=None, wind_speed=None, wind_deg=None, rain_3h=None, clouds_all=None, weather_id=None, weather_main=None, weather_description=None)

At this point I am not sure on what to do?

Do I go ahead and create temporary columns to manage this and use a regex replacement to fill in these values ?

Or is there any better way I can import the data and manage these mangled values either in sqoop or in pyspark ?

0 Answers
Related