Skip to content

[SPARK-14409][ML][WIP] Add RankingEvaluator - #16618

Closed
daniloascione wants to merge 7 commits into
apache:masterfrom
daniloascione:SPARK-14409
Closed

[SPARK-14409][ML][WIP] Add RankingEvaluator#16618
daniloascione wants to merge 7 commits into
apache:masterfrom
daniloascione:SPARK-14409

Conversation

@daniloascione

@daniloascionedaniloascione commented Jan 17, 2017

Copy link
Copy Markdown

What changes were proposed in this pull request?

This patch adds the implementation of a Dataframe api based RankingEvaluator to ML (ml.evaluation)

How was this patch tested?

Additional test case has been added.

@AmplabJenkins

Copy link
Copy Markdown

Can one of the admins verify this patch?

@daniloascione

Copy link
Copy Markdown
Author

Please consider this PR WIP. Discussion in JIRA https://issues.apache.org/jira/browse/SPARK-14409

@HyukjinKwon

Copy link
Copy Markdown
Member

Could you add [WIP] in the title if it is WIP?

@daniloascionedaniloascione changed the title [SPARK-14409][ML] Add RankingEvaluator[WIP][SPARK-14409][ML] Add RankingEvaluatorJan 18, 2017
@daniloascionedaniloascione changed the title [WIP][SPARK-14409][ML] Add RankingEvaluator[SPARK-14409][ML][WIP] Add RankingEvaluatorMar 6, 2017
@daniloascione

Copy link
Copy Markdown
Author

I rewrote the ranking metrics from the mllib package as UDFs (as suggested here) with minimum changes to the logic.

@MLnick

Copy link
Copy Markdown
Contributor

The basic direction looks right - I won't have time to review immediately. Spark 2.2 QA code freeze will happen shortly so this will wait until 2.3 dev cycle starts

@ebernhardsonebernhardson left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I like where this is going, it could be quite useful for taking advantage of CrossValidation when doing pairwise and listwise ranking in xgboost4j-spark.

val predictionAndLabels: DataFrame = dataset
.join(topAtk, Seq($(queryCol)), "outer")
.withColumn("topAtk", coalesce(col("topAtk"), mapToEmptyArray_()))
.select($(labelCol), "topAtk")

@ebernhardsonebernhardsonApr 26, 2017

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Don't we also need to run an aggregation on the label column, roughly the same as the previous aggregation but using labelCol as the sort instead of predictionCol?

Currently this generates a row per prediction, when ranking tasks should have a row per query. I think the aggregation should be run twice, then those two aggregations should be joined together on queryCol. That would result in a dataset containing (labels of top k predictions, top k actual labels)

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes, I agree. This is currently done in the previous step, when the topAtk Dataframe is calculated (line 101).

Unfortunately this is not compatible with RankingMetrics, which expects the format of predictionAndLabels as input. I didn't want to change RankingMetrics in this same PR.
So the predictionAndLabels DataFrame is calculated to use the same RankingMetrics from mllib (well, it is now UDFs based, but I didn't touched its logic).

var i = 0
while (i < n) {
val gain = 1.0 / math.log(i + 2)
if (i < predicted.length && actualSet.contains(predicted(i))) {

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This doesn't seem right, there is no overlap between the calculation of dcg and max_dcg. The question asked here should be if the label at predicted(i) is "good". When treating the labels as binary relevant/not relevant I suppose that might use a threshold, but better would be to move away from a binary dcg and use the full equation from the docblock. I understand though that you are not looking to make major updates to the code from mllib, so it would probably be reasonable for someone to fix this in a followup.

@daniloascionedaniloascioneApr 27, 2017

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes, this should be fixed in another PR to keep changes isolated. FYI, the original JIRA for this is here.

@MLnick

MLnick commented Jun 26, 2017

Copy link
Copy Markdown
Contributor

@daniloascione are you able to update this? I'd like to target for 2.3.

But, can we do the following:

  1. Port over ranking metrics (udfs) to ml with the RankingEvaluator as you've done, but excluding MPR
  2. Also don't make any logic changes for the metrics calculation

Let's focus on porting things over and getting the API right. Then create follow up tickets for additional metrics (MPR and any others) as well as looking into correcting the logic and/or naming of the existing metrics.

If you are unable to take it up again, I can help.

Thanks!

@MLnick

Copy link
Copy Markdown
Contributor

@daniloascione any update?

@HyukjinKwonHyukjinKwon mentioned this pull request Jul 31, 2017
@asfgitasfgit closed this in 3a45c7fAug 5, 2017
@Kornel

Kornel commented Oct 4, 2017

Copy link
Copy Markdown

@MLnick I'm wondering what's the status of this issue: seems closed, have you any plans on picking it up again?

I might pick it up, but I'm not sure what's left: move from package mllib to ml and maybe a python API? Or fixes to ndcg as well?

@acompa

Copy link
Copy Markdown

I'm also curious about this @MLnick. Seems like there was a lot of movement earlier this year, but this PR has gotten stale.

I can also contribute if @Kornel cannot for any reason.

}
}, DoubleType)

val R_prime = predictionAndObservations.count()

@bantmenbantmenOct 3, 2019

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Shouldn't this be a sum instead of count?
(I know this is old/closed but other people might be referring to this code)

zifeif2 pushed a commit to zifeif2/spark that referenced this pull request Nov 22, 2025
## What changes were proposed in this pull request?
This PR proposes to close stale PRs, mostly the same instances with apache#18017Closesapache#14085 - [SPARK-16408][SQL] SparkSQL Added file get Exception: is a directory …
Closesapache#14239 - [SPARK-16593] [CORE] [WIP] Provide a pre-fetch mechanism to accelerate shuffle stage.
Closesapache#14567 - [SPARK-16992][PYSPARK] Python Pep8 formatting and import reorganisation
Closesapache#14579 - [SPARK-16921][PYSPARK] RDD/DataFrame persist()/cache() should return Python context managers
Closesapache#14601 - [SPARK-13979][Core] Killed executor is re spawned without AWS key…
Closesapache#14830 - [SPARK-16992][PYSPARK][DOCS] import sort and autopep8 on Pyspark examples
Closesapache#14963 - [SPARK-16992][PYSPARK] Virtualenv for Pylint and pep8 in lint-python
Closesapache#15227 - [SPARK-17655][SQL]Remove unused variables declarations and definations in a WholeStageCodeGened stage
Closesapache#15240 - [SPARK-17556] [CORE] [SQL] Executor side broadcast for broadcast joins
Closesapache#15405 - [SPARK-15917][CORE] Added support for number of executors in Standalone [WIP]
Closesapache#16099 - [SPARK-18665][SQL] set statement state to "ERROR" after user cancel job
Closesapache#16445 - [SPARK-19043][SQL]Make SparkSQLSessionManager more configurable
Closesapache#16618 - [SPARK-14409][ML][WIP] Add RankingEvaluator
Closesapache#16766 - [SPARK-19426][SQL] Custom coalesce for Dataset
Closesapache#16832 - [SPARK-19490][SQL] ignore case sensitivity when filtering hive partition columns
Closesapache#17052 - [SPARK-19690][SS] Join a streaming DataFrame with a batch DataFrame which has an aggregation may not work
Closesapache#17267 - [SPARK-19926][PYSPARK] Make pyspark exception more user-friendly
Closesapache#17371 - [SPARK-19903][PYSPARK][SS] window operator miss the `watermark` metadata of time column
Closesapache#17401 - [SPARK-18364][YARN] Expose metrics for YarnShuffleService
Closesapache#17519 - [SPARK-15352][Doc] follow-up: add configuration docs for topology-aware block replication
Closesapache#17530 - [SPARK-5158] Access kerberized HDFS from Spark standalone
Closesapache#17854 - [SPARK-20564][Deploy] Reduce massive executor failures when executor count is large (>2000)
Closesapache#17979 - [SPARK-19320][MESOS][WIP]allow specifying a hard limit on number of gpus required in each spark executor when running on mesos
Closesapache#18127 - [SPARK-6628][SQL][Branch-2.1] Fix ClassCastException when executing sql statement 'insert into' on hbase table
Closesapache#18236 - [SPARK-21015] Check field name is not null and empty in GenericRowWit…
Closesapache#18269 - [SPARK-21056][SQL] Use at most one spark job to list files in InMemoryFileIndex
Closesapache#18328 - [SPARK-21121][SQL] Support changing storage level via the spark.sql.inMemoryColumnarStorage.level variable
Closesapache#18354 - [SPARK-18016][SQL][CATALYST][BRANCH-2.1] Code Generation: Constant Pool Limit - Class Splitting
Closesapache#18383 - [SPARK-21167][SS] Set kafka clientId while fetch messages
Closesapache#18414 - [SPARK-21169] [core] Make sure to update application status to RUNNING if executors are accepted and RUNNING after recovery
Closesapache#18432 - resolve com.esotericsoftware.kryo.KryoException
Closesapache#18490 - [SPARK-21269][Core][WIP] Fix FetchFailedException when enable maxReqSizeShuffleToMem and KryoSerializer
Closesapache#18585 - SPARK-21359
Closesapache#18609 - Spark SQL merge small files to big files Update InsertIntoHiveTable.scala
Added:
Closesapache#18308 - [SPARK-21099][Spark Core] INFO Log Message Using Incorrect Executor I…
Closesapache#18599 - [SPARK-21372] spark writes one log file even I set the number of spark_rotate_log to 0
Closesapache#18619 - [SPARK-21397][BUILD]Maven shade plugin adding dependency-reduced-pom.xml to …
Closesapache#18667 - Fix the simpleString used in error messages
Closesapache#18782 - Branch 2.1
Added:
Closesapache#17694 - [SPARK-12717][PYSPARK] Resolving race condition with pyspark broadcasts when using multiple threads
Added:
Closesapache#16456 - [SPARK-18994] clean up the local directories for application in future by annother thread
Closesapache#18683 - [SPARK-21474][CORE] Make number of parallel fetches from a reducer configurable
Closesapache#18690 - [SPARK-21334][CORE] Add metrics reporting service to External Shuffle Server
Added:
Closesapache#18827 - Merge pull request 1 from apache/master
## How was this patch tested?
N/A
Author: hyukjinkwon <gurwls223@gmail.com>
Closesapache#18780 from HyukjinKwon/close-prs.
Sign up for freeto 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.

8 participants

@daniloascione@AmplabJenkins@HyukjinKwon@MLnick@Kornel@acompa@ebernhardson@bantmen