Skip to content

Connect apply labels - #1408

Open
Aditya Pareek (LOGANBLUE1) wants to merge 5 commits into
zinggAI:spark-connectfrom
LOGANBLUE1:connect-apply-labels
Open

Aditya Pareek (LOGANBLUE1) wants to merge 5 commits into
zinggAI:spark-connectfrom
LOGANBLUE1:connect-apply-labels

Conversation

@LOGANBLUE1

Copy link
Copy Markdown
Contributor

No description provided.

…ronment

The spark-connect proto and server-plugin modules each carried their own
copy of the four Spark profiles purely to set <protobuf.version>. Move that
property onto the root profiles and manage protobuf-java in the root
<dependencyManagement>, so the child modules can declare the dependency
without a version -- Maven cannot interpolate an inherited property into a
child's dependency.version at model-build time. Also aligns both parent
references with the ${zingg.version} convention already used by assembly.

zc.sh was pinned to one machine's paths, Spark 3.5.1 and protobuf 3.23.4.
It now derives the Spark and Scala versions from the spark-core jar under
SPARK_HOME and picks the matching protobuf, venv and model dir from that,
so switching Spark lines is just a different SPARK_HOME. Adds a preflight
check, an `env` subcommand, a Java 17+ check for the 4.x line (verifying
each candidate, since `java_home -v 17` returns a wrong JDK rather than
failing), a 4g driver default for the match phase, and only passes
--packages spark-connect when the distribution does not bundle it.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
SparkConnectClient.to_table gained a required `observations` argument in
pyspark 4.x and returns a 3-tuple there (table, schema, execution_info)
against 2 on 3.5, so getUnmarkedPairs() failed on the 4.x line. Probe the
signature rather than branching on a version string.

Also regenerates zingg_command_pb2.py with the protoc the 4.2 profile
pins. Note this generated file is profile-dependent: it now asserts a
Python protobuf runtime >= 6.33.5, while pyproject.toml still declares
protobuf>=4.21.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
… pipeline

The Spark Connect label loop wrote the marked pipe from the client: it read
the unmarked parquet over Connect, joined the client's decisions onto it and
appended to the marked folder. That duplicated LabellerUtil.postProcessLabel
on the client, so any front end that got the row reshaping subtly wrong wrote
subtly wrong training data, and each new front end would have to repeat it.

Introduce ILabelCollector as the seam between Labeller's pipeline and whatever
supplies the labels. Labeller.execute() now reads, preprocesses, calls the
collector, postprocesses and writes; only the collect step varies --
CliLabelCollector prompts a terminal, MapLabelCollector replays decisions made
elsewhere. Labeller.applyLabels(Map) installs the latter and runs the same
execute(), exposed on IZingg/Client (UnsupportedOperationException from
ZinggBaseCommon for the non-labelling phases).

ZinggCommand gains a `labels` map (cluster id -> 1/0/2). A labelling command
carrying labels is recorded via Client.applyLabels; one without still errors,
since the server has no terminal to prompt at. writeMarkedPairs() on the
Python client now sends that command instead of manipulating DataFrames, and
write_marked_pairs/_resolve_session are gone.

ZinggRelationPlugin builds its executor through SparkClient like the command
plugin does, rather than instantiating SparkLabeller/SparkFindAndLabeller
itself -- and no longer runs the finder during plan transformation, which
Spark may invoke more than once; findAndLabel callers run findTrainingData as
its own command first.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant