Apache Spark — The answer to a slow job is in the plan and the event log
Read a partner's CSV with broken rows in three different modes
Goal
You read the 6,000-row order file a partner sent with a schema you define, confirm with numbers how each of the three modes PERMISSIVE, DROPMALFORMED and FAILFAST handles the broken rows, and then add business rules and split the data into a clean table and a quarantine table.
Why it matters
A CSV has no types. If "two" slips into one column even once, inferSchema falls back to a string for that whole column, and this shows up not as an error but as a silent type change. If a column that was an integer until yesterday becomes a string today, every total after it breaks. That is why a pipeline does not infer the schema but defines it when reading.
Once you define the schema, you then have to decide what to do with rows that do not fit it. Spark's CSV reading has three modes. PERMISSIVE keeps the broken row and puts the original text in the corrupt-record column, DROPMALFORMED discards it, and FAILFAST stops. Which one is right is decided by the business: for a table that counts money, silently discarding is the most dangerous.
There is one trap. The CSV parser parses only the columns it needs. An action that needs no column at all, such as count(), does not parse the rows, so if you read with DROPMALFORMED and count, you count even the broken rows, and FAILFAST does not stop either. If you did the validation with count(), you validated nothing.
Steps
- Write the seven columns of the original as DDL and save them to /root/spk/schema/schema.ddl (
qtyis INT,order_tsis TIMESTAMP, the rest are STRING). - Create /root/spk/schema/infer.py under the application name
spk-schema-infer, read withinferSchema=True, and write thesimpleString()of the inferred schema to /root/spk/schema/out/inferred.txt. - Create /root/spk/schema/load.py under the application name
spk-schema-load, add_corrupt STRINGto the end of the step 1 schema, read with PERMISSIVE, and write everything as Parquet to /root/spk/schema/out/permissive. - Create /root/spk/schema/dropmal.py under the application name
spk-schema-drop, read with DROPMALFORMED, and write the value fromcount()and the value counted again after writing as Parquet to /root/spk/schema/out/dropped into /root/spk/schema/out/drop_count.json as{"count_only": 정수, "written": 정수}(both values are integers). - Create /root/spk/schema/failfast.py under the application name
spk-schema-failfast, make it write what it read with FAILFAST as Parquet, and write the name of the parsing error condition that caused the failure to the first line of /root/spk/schema/out/failfast.txt. - Create /root/spk/schema/width.py under the application name
spk-schema-width, read the original line by line (excluding the header row), and write the number of lines per number of cells when split on commas to /root/spk/schema/out/width.json as{"칸 수": 줄 수}(the key is the cell count as a string and the value is the line count). - Create /root/spk/schema/clean.py under the application name
spk-schema-clean, and among the rows that were parsed, write only those withqtyof 1 or more to /root/spk/schema/out/clean, and the rest, with areasoncolumn added (parsefor parse failures,qtyfor quantity problems), to /root/spk/schema/out/quarantine, both as Parquet. - In /root/spk/schema/report.md, write three sections:
## 추론이 틀린 곳,## 세 가지 모드and## 격리한 줄(use exactly these Korean headings in this order; they mean "Where inference is wrong", "The three modes" and "Quarantined rows"). In the second section, put the number of corrupt rows and the number of rows written with DROPMALFORMED; in the third, put the number of clean rows and the number of quarantined rows, as numbers.
Notes
- Original:
/data/shop/orders_dirty.csv(1 header row + 6,000 rows). The time format isyyyy-MM-dd HH:mm:ss, so tell Spark withtimestampFormat. - The corrupt-record column must be put directly into the schema. Its name must be the same as the
columnNameOfCorruptRecordoption. - A query that selects only the corrupt-record column may be blocked because it would have to re-parse the original.
cache()what you read once, then filter. - The exception from FAILFAST may not return the condition name directly in Python. The condition name is inside the square brackets
[...]in the message. - Common mistakes: checking the effect of a mode with
count(); not putting the corrupt-record column in the schema; sending an empty quantity (null) to the clean rows. - Official documentation: CSV Files · Data Types · Datetime Patterns · Error Conditions
Define the schema up front
Write the seven columns, in the order of the header row of the original /data/shop/orders_dirty.csv, as a one-line DDL string and save it to /root/spk/schema/schema.ddl. qty is INT, order_ts is TIMESTAMP, and the rest are STRING.
DDL has the shape 이름 형, 이름 형, ... (name, type, name, type, and so on). Look at the header row with head -1 and follow its order exactly: a CSV matches columns by position, not by name.
See what inference misses
Create /root/spk/schema/infer.py under the application name spk-schema-infer, read the original with inferSchema=True and header=True, and write df.schema.simpleString() to /root/spk/schema/out/inferred.txt.
Inference is a job that scans the data one more time. If you look at the job count in the event log, it is higher than when you defined the schema. In the result, look at what types qty and order_ts were given, and grep the original to see why.
PERMISSIVE: keep the original text instead of discarding
Create /root/spk/schema/load.py under the application name spk-schema-load, read with a schema made by adding _corrupt STRING to the end of the step 1 DDL, with mode=PERMISSIVE, columnNameOfCorruptRecord=_corrupt and timestampFormat=yyyy-MM-dd HH:mm:ss, and write everything as Parquet to /root/spk/schema/out/permissive.
PERMISSIVE does not change the number of rows. All 6,000 rows remain, and for a row that does not match the type or has a different number of cells, the original text goes into _corrupt and the affected cells become null. The grader re-reads the original, counts the broken rows and compares them with your _corrupt.
DROPMALFORMED: count() lies
Create /root/spk/schema/dropmal.py under the application name spk-schema-drop, read with the step 1 schema and mode=DROPMALFORMED, and write (a) the value from count() and (b) the value counted after writing Parquet to /root/spk/schema/out/dropped and reading it back, into /root/spk/schema/out/drop_count.json as {"count_only": 정수, "written": 정수} (both values are integers).
The two numbers should come out different. count() needs no column, so the CSV parser does not parse the rows, and if it does not parse, it does not know whether a row is broken. It discards rows only when you perform an action that has to use all the columns.
FAILFAST: make it stop and read the reason
Create /root/spk/schema/failfast.py under the application name spk-schema-failfast, make it write what it read with the step 1 schema and mode=FAILFAST as Parquet (any path), catch the failure, and write the name of the parsing-failure error condition (the one that starts with MALFORMED_…) to the first line of /root/spk/schema/out/failfast.txt.
A write parses all the columns, so the job fails at the first broken row. The exception is wrapped in layers: the outer one is a file read failure, and the parse failure is on the cause side. Extract every name inside square brackets from the message. The grader also checks whether that application actually had a failed job.
Count the rows with a different number of cells
Create /root/spk/schema/width.py under the application name spk-schema-width, read the original line by line with spark.read.text, count the lines other than the header row by the number of cells when split on commas, and write them to /root/spk/schema/out/width.json as {"칸 수": 줄 수} (the key is a string: the cell count, and the value is the line count).
F.size(F.split("value", ",")) is the number of cells. This original has no commas inside quotes, so simple splitting is fine (with a CSV that has quotes, this method would be wrong). Compare how much of the corrupt rows from step 3 is taken up by rows that do not have 7 cells.
Business rules too: a clean table and a quarantine table
Create /root/spk/schema/clean.py under the application name spk-schema-clean, read with PERMISSIVE, and write as Parquet only the rows that were parsed (_corrupt is null) and have qty of 1 or more to /root/spk/schema/out/clean (without the corrupt-record column), and the rest, with a reason column (parse or qty) added, to /root/spk/schema/out/quarantine.
There are rows whose format is right but that are wrong for the business: negative quantities and empty quantities. The parser does not see these as broken rows, so you have to write the rule yourself. If you keep the reason as a column, it is immediately clear what to ask the partner to fix when you send the rows back.
Record what you discarded and why
In /root/spk/schema/report.md, write three sections: ## 추론이 틀린 곳, ## 세 가지 모드 and ## 격리한 줄 (use exactly these Korean headings in this order; they mean "Where inference is wrong", "The three modes" and "Quarantined rows"). In the second section, put the number of corrupt rows from step 3 and the number of rows written and counted in step 4; in the third, put the number of clean rows and the number of quarantined rows from step 7, as numbers.
Re-count the numbers from your result files and copy them over. Think of it as an email to the partner: how many rows you received, how many rows broke and why, and how many rows you are sending back because of business rules.