JSON is where most data pipelines begin: application logs, webhook payloads, API exports and event streams all arrive as JSON. Spark reads it with one line, spark.read.json(path), and that convenience is the problem. The default path reads your data twice to guess a schema, quietly turns bad records into nulls, and makes different guesses on different days. Most JSON incidents in Spark trace back to those three defaults.

This article explains how the JSON source actually works: the two file layouts and why one cannot be split, how schema inference decides types, what each parse mode does with malformed input, the options worth changing, nested data and the VARIANT type added in Spark 4.0, and how to write JSON back out. Option names and defaults are taken from the Spark 4.2.0 documentation. By the end you should be able to turn a fragile read.json into an explicit, observable ingestion step.

JSON Lines and multi-line documents

Spark's JSON source expects JSON Lines by default: one complete JSON object per line, with lines separated by \n, \r\n or \r. This is what most loggers and exporters produce, and it is splittable: a 10 GB uncompressed file can be cut into byte ranges and each task finds the next line boundary and parses from there.

A file that is a single JSON document, such as a pretty-printed array of objects, needs multiLine=true (default false). In that mode each file is parsed as a whole by one task, so one 8 GB file means one task and one executor holding the work, however large the cluster. Compression adds a second constraint: gzip files are never splittable, so a single large .json.gz is also one task even in JSON Lines mode, while bzip2 is splittable but slow to decompress. If you control the producer, ask for JSON Lines in many moderately sized files; if you do not, convert once to Parquet and read that thereafter.

How a JSON read runs

JSON filesJSON Lines or documentsSplit planningsplittable?Schemagiven or inferredInference passextra full readParse tasksJackson per recordGood rowstyped DataFrameMalformed rowsper modeWriteParquet / DeltaVARIANTparse_json (4.0)no schemaschemaraw stringPERMISSIVE keeps malformed text in a corrupt-record column, DROPMALFORMED drops it,FAILFAST stops the job. Inference costs a full extra read unless you sample or supply a schema.
Figure 1. The JSON read path. Without a schema Spark adds a full inference pass; with one it parses directly. Malformed records follow the configured mode; raw strings can instead be kept as VARIANT.

Reading happens in two phases. Planning lists files and creates input splits. Then, if you did not supply a schema, Spark runs a job that parses records and merges their inferred types into one schema; only after that does the real read start. With samplingRatio at its default of 1.0, inference parses every record, so an unspecified schema doubles I/O and CPU on the first action. Each record is parsed by Jackson into the target row format, and anything that does not fit is handled by the parse mode.

Schema inference and explicit schemas

Inference walks every sampled record and widens types as it goes. Integers become long, decimals become double (or decimal with prefersDecimal=true), and a field that is a number in one record and a string in another becomes string. Fields missing from some records become nullable columns, and objects become structs whose field set is the union across records. Under the default settings, dates and timestamps arrive as strings. Fields that are always null become strings and always-empty arrays become arrays of strings, unless dropFieldIfAllNull=true removes them.

Three consequences follow. Sampling, for example samplingRatio=0.01, is cheaper but can miss a rare field or a rare type, and the result can change between runs. The inferred schema depends on the data in today's files, so a producer adding a field or sending one string where it used to send numbers changes your table's types without any code change. And inference cannot see intent: an ID that looks numeric is inferred as long and loses leading zeros.

The fix is to declare the schema in code, usually as a DDL string, and infer only during exploration:

from pyspark.sql import functions as F

EVENT_SCHEMA = '''
  event_id STRING,
  user_id STRING,
  ts TIMESTAMP,
  type STRING,
  amount DECIMAL(12,2),
  device STRUCT<os: STRING, version: STRING>,
  items ARRAY<STRUCT<sku: STRING, qty: INT>>,
  _corrupt_record STRING
'''

events = (spark.read
    .schema(EVENT_SCHEMA)
    .option("mode", "PERMISSIVE")
    .option("timestampFormat", "yyyy-MM-dd'T'HH:mm:ss[.SSS][XXX]")
    .json("s3://raw/events/dt=2026-10-04/"))

# During exploration only: infer from a sample and print a schema you can paste.
print(spark.read.option("samplingRatio", 0.05).json("s3://raw/events/sample/").schema.simpleString())

Parse modes and corrupt records

The mode option decides what happens when a record cannot be parsed or does not fit the schema. PERMISSIVE, the default, sets the fields it cannot read to null and, if the schema contains a string column named by columnNameOfCorruptRecord (default _corrupt_record), stores the raw text there. If you supply a schema without that column, bad records simply become rows of nulls and the evidence is gone. DROPMALFORMED discards bad records silently. FAILFAST throws on the first one, which is right when bad input should stop a pipeline rather than corrupt a table.

Since Spark 2.3, a query that references only the corrupt-record column of a raw JSON read is rejected, because column pruning would let Spark skip parsing the other fields and so it could not know which records are malformed. Cache or materialise the parsed DataFrame first, then split it:

parsed = events.cache()
bad = parsed.filter(F.col("_corrupt_record").isNotNull())
good = parsed.filter(F.col("_corrupt_record").isNull()).drop("_corrupt_record")

bad_count = bad.count()
total = parsed.count()
if total and bad_count / total > 0.001:          # tune the budget to your source
    raise RuntimeError(f"{bad_count} of {total} records malformed")

bad.write.mode("append").json("s3://quarantine/events/dt=2026-10-04/")
good.write.mode("overwrite").parquet("s3://bronze/events/dt=2026-10-04/")

This is the pattern to aim for: keep bad records in a quarantine location, measure them, and fail the run when the rate passes a budget. Silent nulls are the worst outcome, because they look like real data in every downstream aggregate.

Options worth knowing

OptionDefaultWhen to change it
multiLinefalseinput is whole JSON documents or arrays, not JSON Lines
samplingRatio1.0only when inferring; lower it for exploration
primitivesAsStringfalseland everything as strings and cast later
prefersDecimalfalsemoney and exact quantities
allowCommentsfalsehand-written config files
allowSingleQuotestrueset false to be strict about standard JSON
allowNonNumericNumberstrueset false to reject NaN and INF tokens
dateFormat, timestampFormatISO-8601 patternsproducers using other layouts
encodingauto-detected for multiLine readsUTF-16 or UTF-32 sources
lineSepany of \r, \r\n, \nunusual record separators
dropFieldIfAllNullfalseinference produces useless all-null columns

Two defaults are permissive in a way that surprises people: single-quoted strings and NaN/INF tokens are accepted unless you turn them off. If a downstream system requires strict JSON, validate strictly at ingestion rather than discovering the difference later.

Nested data, from_json and VARIANT

JSON data is rarely flat. Struct fields are addressed with dots, arrays are expanded into rows with explode (or explode_outer to keep rows whose array is null or empty), and higher-order functions such as transform and filter work inside arrays without exploding them.

JSON often arrives as a string inside another format: a Kafka value, a database column, a field in a CSV. from_json parses such a column with an explicit schema, schema_of_json infers a schema from a sample literal, and to_json serialises a struct back to text. get_json_object extracts one path from a string without a schema, but it re-parses the string for every call, so use from_json once when you need several fields.

lines = (events
    .select("event_id", "ts", F.col("device.os").alias("os"),
            F.explode_outer("items").alias("item"))
    .select("event_id", "ts", "os", "item.sku", "item.qty"))

payload_schema = "order_id STRING, total DECIMAL(12,2), tags ARRAY<STRING>"
orders = raw.select(F.from_json("value", payload_schema).alias("o")).select("o.*")

Spark 4.0 added a third option for data whose shape is genuinely unknown or varies by record: the VARIANT type. parse_json turns a JSON string into a variant that stores the structure in a binary encoding, and variant_get extracts a path and casts it, returning null when the path is absent. That lets you land data without fixing a schema and decide on columns later, without re-parsing text on every query.

raw = spark.read.text("s3://raw/webhooks/")            # one JSON document per line
v = raw.select(F.parse_json("value").alias("doc"))
v.select(F.variant_get("doc", "$.customer.id", "string").alias("customer_id"),
         F.variant_get("doc", "$.amount", "decimal(12,2)").alias("amount"))

Writing JSON

df.write.json(path) writes JSON Lines, one file per partition. Write options include compression (none, bzip2, gzip, lz4, snappy or deflate), ignoreNullFields (from spark.sql.jsonGenerator.ignoreNullFields, which drops null fields from output objects) and the date and timestamp formats. Use partitionBy for directory layout and control the file count with repartition or coalesce before the write; see coalesce and repartition.

Write JSON when the consumer needs it: an external API, a search index loader, a partner export. For data Spark will read again, write Parquet or a table format instead; columnar files are smaller, typed and support column pruning and predicate pushdown, which JSON cannot. Spark Parquet read and write covers that side.

Worked example: a slow clickstream job

A team ingests clickstream exports: each hour, one gzip file of about 6 GB compressed, JSON Lines inside. The job uses spark.read.json with no schema and takes 50 minutes on a 40-executor cluster, mostly idle.

The diagnosis starts in the Spark UI. The inference stage and the read stage each show one task per file, because gzip cannot be split, so the cluster runs one long task twice. The fixes, in order: supply a DDL schema, which removes the inference pass and halves the work; ask the producer to write many smaller files or a splittable format, or, if that is impossible, read once and immediately repartition(400) so everything after the decompression runs in parallel; and write the result to Parquet partitioned by date and hour, so downstream jobs never touch the JSON again.

Adding the corrupt-record column during the change exposes the next issue: 0.3% of records have a truncated last line from an exporter bug, which had been turning into rows of nulls and inflating the anonymous visitor count. The job now quarantines them and fails above 0.1%, and the exporter fix ships the following week. Total run time after the changes is dominated by single-threaded gzip decompression, which is the remaining argument for changing the producer.

Failure modes

  • Schema drift. A producer change flips a column's type or adds a field and the inferred schema follows it. Declare schemas and compare incoming data against them in a check.
  • Silent nulls. A schema without the corrupt-record column turns bad input into null rows. Always include it in PERMISSIVE mode, and count it.
  • One giant task. multiLine documents or large gzip files run as one task per file. Look for a stage with few, long tasks.
  • Timestamp mismatch. A format that does not match the data yields nulls, not errors, in PERMISSIVE mode. Test the format on real samples and check null rates.
  • Case-colliding keys. Objects with keys that differ only in case, such as id and ID, can make column references ambiguous because Spark resolves names case-insensitively by default.
  • Small files. Millions of tiny JSON files make planning slow and tasks inefficient. Compact them at landing, and see Spark partitioning.

Trade-offs

Explicit schemas give stable types and faster reads but must be maintained as producers change. Inference adapts automatically but makes your table's contract depend on the data. VARIANT sits between: it accepts any shape cheaply and defers decisions, at the cost of weaker typing until you extract columns. A common layout is to land raw JSON as VARIANT or strings in a bronze table, then promote stable fields to typed columns in silver.

Parse modes trade completeness against correctness: PERMISSIVE with quarantine is the default to choose for ingestion, FAILFAST for contracts that must hold, and DROPMALFORMED almost never, since it hides the problem it exists to handle.

What to do next

  1. Find every read.json without a .schema(...) in your jobs and list them.
  2. Generate a schema from a sample once, review the types (IDs as strings, money as decimal) and commit it.
  3. Add _corrupt_record to each schema, quarantine bad rows and fail above a budget.
  4. Check input files for gzip or multiLine and fix the parallelism in the UI.
  5. Set allowSingleQuotes and allowNonNumericNumbers to false where strict JSON matters.
  6. Convert JSON to Parquet or a table format at the first hop; see Delta Lake on Spark.
  7. On Spark 4.0 or later, try parse_json for sources whose shape you cannot pin down.
Key takeaway: Treat spark.read.json as an ingestion boundary, not a convenience. Prefer JSON Lines in splittable files, declare schemas instead of inferring them, keep the corrupt-record column and quarantine bad rows against a budget, and be explicit about timestamp formats and strictness. Use from_json for JSON inside other formats and VARIANT on Spark 4.0 or later for shapes you cannot fix, and convert to Parquet or a table format as early as possible.