Koalas was the Databricks open-source project that put a pandas-shaped API on top of Spark DataFrames. In Spark 3.2 it was merged into PySpark itself as pyspark.pandas, now documented as the pandas API on Spark. The standalone Koalas package stayed on older Spark releases and is no longer where development happens, so any job still importing databricks.koalas is pinned to an old runtime and blocks every platform upgrade around it.

The migration looks like a find-and-replace, and the first step really is one. The trouble is that most teams do the rename and the Spark upgrade together, and the upgrade brings behaviour changes that no rename can catch: a different default index, a different SQL templating rule, removed methods, and in Spark 4.0 an exception when ANSI mode is on. This article splits the work into two hops, gives the exact renames from the official migration notes, lists the changes per release, and shows a codemod and a parity-test harness so you can prove each hop before shipping it. How pandas-on-Spark works internally is covered in the Spark pandas API article; this page is only about getting there safely.

Advertisement

What actually changed between Koalas and pyspark.pandas

The pandas API on Spark is the Koalas code base moved into the Spark repository, so the core idea is unchanged: a pandas-on-Spark DataFrame wraps a Spark DataFrame plus an internal frame that tracks the index and column labels, and every pandas-style call compiles to a Spark plan. What changed is the surface around it. The PySpark documentation lists the migration notes, and there are only a handful:

Koalaspandas API on SparkStatus
import databricks.koalas as ksimport pyspark.pandas as psPackage renamed
kdf.koalas.apply_batch(...)psdf.pandas_on_spark.apply_batch(...)Accessor renamed; .koalas removed in Spark 4.0
sdf.to_koalas()sdf.pandas_api()Removed in Spark 4.0
sdf.to_pandas_on_spark()sdf.pandas_api()Interim 3.2 name; removed in Spark 4.0
databricks.koalas.__version__pyspark.__version__Removed; the version is Spark's

Two consequences follow. First, the version of your pandas API is now the version of Spark: you cannot upgrade one without the other, and your pandas, NumPy and PyArrow versions must satisfy Spark's floors. Second, the old names did not vanish at once. They were kept as deprecated aliases through the 3.x line, which is why code that half-migrated years ago still runs on 3.5 and then fails on 4.0.

Option names did not change: compute.default_index_type, compute.ops_on_diff_frames and the rest are set with ps.set_option exactly as with ks.set_option. Their defaults are another matter, covered below.

Plan two hops, not one

Treat the migration as two separate releases, each with its own test gate. Hop 1 moves from Koalas to pyspark.pandas on a Spark 3.x runtime, ideally the newest 3.5 release you can run. Hop 2 moves that code to Spark 4.x.

Two hops, each with its own breakage: rename first, then upgrade Spark one major step at a timeKoalas 1.xdatabricks.koalas on Spark 3.0/3.1pyspark.pandasSpark 3.2 - 3.5pyspark.pandasSpark 4.xhop 1hop 2Hop 1 changesimport path.koalas -> .pandas_on_sparkto_koalas -> pandas_api__version__ removedInside 3.x3.3: sql() formatter3.3: drop() by index3.4: pandas 2.0 deprecations3.3: pandas >= 1.0.5Hop 2 changesdeprecated APIs removedANSI mode raisesops_on_diff_frames = Truepandas 2 / PyArrow 11Parity tests gate every hopsame fixtures, old vs new output, sorted by keyRename-only diffs are safe to review; behaviour changes need tests, not eyeballs.
Hop 1 is mostly mechanical renames plus the 3.3 and 3.4 behaviour changes; hop 2 removes deprecated APIs, raises dependency floors and adds the ANSI-mode check. Both hops are gated by the same parity tests.

The changes inside the 3.x line matter if you jump from Spark 3.1 straight to 3.5. These are the ones from the PySpark upgrade guide most likely to change results in Koalas-era code:

ReleaseChangeWhy it bites migrated code
3.3ps.sql follows the standard Python string formatterQueries that relied on {name} being filled from the caller's local variables must pass values explicitly; PYSPARK_PANDAS_SQL_LEGACY=1 restores the old behaviour temporarily
3.3DataFrame.drop drops by index by defaultA bare drop(labels) meant for columns now targets rows; write columns= explicitly
3.3Minimum pandas raised to 1.0.5Old pinned environments fail at import
3.4Follows pandas 2.0 deprecations; ps.sql gains named argsDeprecation warnings appear on 3.4 and become errors on 4.0
4.0Deprecated APIs removed; pandas >= 2.0, NumPy >= 1.21, PyArrow >= 11.0Covered in its own section below
Advertisement

Inventory every Koalas touchpoint first

Before changing anything, find every place Koalas is used, including the indirect ones: notebooks, job definitions that install the koalas wheel, requirements files, and code that calls Koalas-only accessors without importing the package directly, for example a function that receives a frame and calls .koalas. on it. A plain search is enough to start:

# source, notebooks and dependency pins
grep -rn --include='*.py' --include='*.ipynb' -E \
  'databricks\.koalas|import koalas|\.koalas\.|to_koalas\(|to_pandas_on_spark\(|ks\.__version__' .
grep -rn -E '^koalas([=<>~ ]|$)' requirements*.txt setup.cfg pyproject.toml 2>/dev/null

# methods removed in Spark 4.0 that Koalas-era code commonly uses
grep -rn --include='*.py' -E '\.(iteritems|append|mad|get_dtype_counts|is_monotonic)\(' .

Turn the output into a tracking sheet with one row per job: owner, Spark version today, which accessors it uses, whether it has tests, and whether its output feeds something that is compared numerically downstream, such as a model or a finance report. The last column decides how strict the parity test must be.

The mechanical rewrite: a small codemod

The renames are safe to automate because they are one-to-one. Keep the codemod narrow: it should only touch the names in the table above, so that its diff can be reviewed quickly and every other change is a deliberate human edit.

import pathlib
import re
import sys

RULES = [
    (re.compile(r"import databricks\.koalas as ks\b"), "import pyspark.pandas as ps"),
    (re.compile(r"from databricks import koalas as ks\b"), "import pyspark.pandas as ps"),
    (re.compile(r"\bdatabricks\.koalas\b"), "pyspark.pandas"),
    (re.compile(r"\bks\.__version__\b"), "pyspark.__version__"),
    (re.compile(r"\bks\."), "ps."),
    (re.compile(r"\.koalas\."), ".pandas_on_spark."),
    (re.compile(r"\.to_koalas\("), ".pandas_api("),
    (re.compile(r"\.to_pandas_on_spark\("), ".pandas_api("),
]

def rewrite(path: pathlib.Path) -> int:
    src = path.read_text(encoding="utf-8")
    out, total = src, 0
    for pattern, repl in RULES:
        out, n = pattern.subn(repl, out)
        total += n
    if total:
        path.write_text(out, encoding="utf-8")
    return total

if __name__ == "__main__":
    for p in sys.argv[1:]:
        n = rewrite(pathlib.Path(p))
        print(f"{p}: {n} rewrites")

Run it file by file, read the diff, then run the import-time check: import every rewritten module in a Spark 3.5 environment. The ks. rule is the one to review carefully, because a different object might be called ks in some files. If pyspark.__version__ is now referenced, add import pyspark.

Semantic changes a rename cannot fix

These changes keep the code running and change what it returns. Each needs a human decision.

  • Default index type. Every frame created from Spark data gets a generated index. The pandas API on Spark defaults compute.default_index_type to distributed-sequence; Koalas documentation listed sequence as its default, so check the release you are leaving. Both produce 0, 1, 2 and so on, but which row gets which number is not guaranteed to be the same across runs with distributed-sequence. Code that joins two frames on a generated index, or that uses iloc to take the first rows of unsorted data, can quietly pair the wrong rows. Fix it by setting a real key with set_index or by sorting before positional access. The Spark pandas API article explains the index types in detail.
  • Operations between different frames. Adding a column from one frame to another needs a join on the index. Koalas and Spark 3.x raised an error unless compute.ops_on_diff_frames was true. Spark 4.0 turns it on by default, so code that used to fail loudly now runs an expensive join, and on a non-deterministic index it can align the wrong rows.
  • SQL templating. From 3.3, ps.sql('SELECT * FROM {t} WHERE x > {lo}', t=psdf, lo=10) is the supported form. Pass frames and values as keyword arguments rather than relying on local-variable capture.
  • Defaults that followed pandas 2. In 4.0, Series.str.replace defaults to regex=False, value_counts names its result count, and datetime attributes such as year are int32.

None of these produce an error on the happy path, which is why parity tests, not code review, are the gate.

Spark 4.0: removals, dependency floors and ANSI mode

Hop 2 is where deprecated code stops running. The upgrade guide removes, among others, DataFrame.append and Series.append (use ps.concat), iteritems (use items), mad, get_dtype_counts, is_monotonic (use is_monotonic_increasing), Int64Index and Float64Index, and the three Koalas-era names in the earlier table. Spark 4.0 also raises its minimums to pandas 2.0.0, NumPy 1.21 and PyArrow 11.0.0, so the Python environment usually moves at the same time.

The most visible change is ANSI mode. Spark 4.0 enables spark.sql.ansi.enabled by default, and the pandas API on Spark raises an exception when it detects ANSI mode, because pandas semantics such as integer division producing nulls do not hold under ANSI rules. You have two documented choices:

from pyspark.sql import SparkSession
import pyspark.pandas as ps

spark = (SparkSession.builder
         .appName("features")
         # Option 1 (recommended for pandas-API jobs): turn ANSI off for this session.
         .config("spark.sql.ansi.enabled", "false")
         .getOrCreate())

# Option 2: keep ANSI on and force the pandas API to run anyway.
# The upgrade guide warns this "can cause unexpected behavior".
# ps.set_option("compute.fail_on_ansi_mode", False)

Prefer option 1 for jobs that are mostly pandas API code, and set it per job rather than cluster-wide so SQL jobs keep ANSI checks. If you run through Spark Connect, test the pandas API path there explicitly, because the client and server are separate processes with separate dependency sets.

A parity-test harness

A parity test runs the same transformation on the old and new stacks with the same input and compares the outputs. Because Koalas and Spark 4 cannot share one Python environment, the practical version runs each side as its own job against a frozen fixture and writes Parquet; a third step compares the two. Sort by a business key before comparing, since row order is not part of the contract.

# compare_outputs.py: run after both the old-stack and new-stack jobs have written Parquet
import pyspark.pandas as ps
from pyspark.pandas.testing import assert_frame_equal

KEYS = ["customer_id", "day"]

old = ps.read_parquet("gs://bucket/parity/features_old").sort_values(KEYS).reset_index(drop=True)
new = ps.read_parquet("gs://bucket/parity/features_new").sort_values(KEYS).reset_index(drop=True)

assert list(old.columns) == list(new.columns), (old.columns, new.columns)
print("rows", len(old), len(new))
assert len(old) == len(new)

# Exact for keys and integers; tolerance for floats that went through different code paths.
assert_frame_equal(old[KEYS], new[KEYS])
assert_frame_equal(old, new, check_exact=False, rtol=1e-9, check_dtype=False)

pyspark.pandas.testing.assert_frame_equal is the replacement the upgrade guide names for the removed assertPandasOnSparkEqual. Use check_dtype=False only for known widening changes such as int32 date parts, and record each exception in the test. For large outputs, compare row counts, per-column null counts and sums on the full data, and run the exact comparison on a sample of keys.

Worked example: a daily feature job

Here is a typical Koalas job: it reads events, derives daily features and joins a lookup table.

import databricks.koalas as ks
ks.set_option("compute.ops_on_diff_frames", True)

ev = ks.read_parquet("/data/events")
daily = ev.groupby(["customer_id", "day"]).agg({"amount": "sum", "event_id": "count"})
daily.columns = ["amount", "n_events"]
daily = daily.reset_index()
tiers = ks.read_csv("/ref/tiers.csv")
daily["tier"] = tiers["tier"]                       # aligned on generated index!
daily = daily.append(ks.read_parquet("/data/backfill"))
daily.to_parquet("/out/features")

The codemod fixes the import and the ks. calls. Review finds two real bugs. The tier assignment aligns two unrelated frames on their generated indexes, which only worked by accident when both happened to be numbered in the same order, so it becomes a keyed merge. And append is gone in Spark 4.0, so it becomes ps.concat. The migrated version:

import pyspark.pandas as ps

ev = ps.read_parquet("/data/events")
daily = (ev.groupby(["customer_id", "day"])
           .agg({"amount": "sum", "event_id": "count"})
           .rename(columns={"event_id": "n_events"})
           .reset_index())
tiers = ps.read_csv("/ref/tiers.csv")[["customer_id", "tier"]]
daily = daily.merge(tiers, on="customer_id", how="left")   # explicit key, no index alignment
daily = ps.concat([daily, ps.read_parquet("/data/backfill")], ignore_index=True)
daily.to_parquet("/out/features")

The parity test then shows a difference in tier for some customers. That difference is the old bug, and it has been feeding a model. Document it, get the data owner's sign-off, and accept the new output as correct. Check the new physical plan with the plan-reading techniques: the merge should appear as one join, not as a join plus a window over the whole data set for index generation.

Rollout, performance and failure modes

Ship each hop behind a switch: run old and new jobs side by side for a few scheduled runs, writing to separate locations, and compare automatically. Watch stage-level metrics as well as correctness; a migration that is correct but doubles shuffle volume is not finished. Python-side costs such as Arrow conversion and UDF overhead are covered in the PySpark performance article.

FailureSymptomFix
Half-migrated accessorsRuns on 3.5 with deprecation warnings, AttributeError on 4.0Search for .koalas. and to_koalas again after the codemod; fail CI on deprecation warnings
Index-aligned assignment between framesPlausible but wrong values; silent on 4.0 because ops_on_diff_frames is onUse keyed merge; set the option to false in tests so accidental alignment fails
ANSI exception at first pandas API callJob fails immediately on Spark 4.0Set spark.sql.ansi.enabled=false for the job
Dependency floorImport errors for pandas, NumPy or PyArrowRebuild the environment against Spark 4.0's minimums and pin it
ps.sql placeholders unresolvedError or wrong literal in the queryPass frames and values as keyword arguments
str.replace pattern now literalReplacements stop matchingAdd regex=True where a pattern is intended

What to do next

  1. Search all repositories, notebooks and job definitions for Koalas imports, accessors, to_koalas and the Spark 4.0 removed methods, and build a per-job tracking sheet.
  2. Pick one representative job, freeze an input fixture, and build the parity harness: old job, new job, keyed comparison with assert_frame_equal.
  3. Run the codemod for hop 1 on Spark 3.5, commit the rename separately, and fix the 3.3 and 3.4 behaviour changes it exposes.
  4. Replace every index-aligned assignment between frames with an explicit merge, and set compute.ops_on_diff_frames to false in tests.
  5. Build a Spark 4.0 environment meeting the pandas 2.0, NumPy 1.21 and PyArrow 11 floors, decide the ANSI setting per job, and remove the deleted APIs.
  6. Run old and new jobs side by side for several runs, compare outputs and stage metrics, then cut over and delete the Koalas dependency.
Key takeaway: Moving from Koalas to the pandas API on Spark has two parts. The renames are small and safe to automate: the import path, the accessor, <code>pandas_api()</code> and the version attribute. The behaviour changes across Spark 3.3, 3.4 and 4.0 are not: the default index, SQL templating, removed methods, pandas 2 defaults, cross-frame operations now on by default, and the ANSI-mode exception. Do it in two hops, gate each with keyed parity tests on frozen inputs, and expect some differences to be old bugs that need a decision rather than a fix.