Explore Courses Blog Tutorials Interview Questions
+3 votes
in Big Data Hadoop & Spark by (900 points)

I'm trying to figure out the best way to get the largest value in a Spark dataframe column.

Consider the following example:

df = spark.createDataFrame([(1., 4.), (2., 5.), (3., 6.)], ["A", "B"])

Which creates:

|  A|  B|

My goal is to find the largest value in column A (by inspection, this is 3.0). Using PySpark, here are four approaches I can think of:

# Method 1: Use describe()
float(df.describe("A").filter("summary = 'max'").select("A").collect()[0].asDict()['A'])

# Method 2: Use SQL
spark.sql("SELECT MAX(A) as maxval FROM df_table").collect()[0].asDict()['maxval']

# Method 3: Use groupby()

# Method 4: Convert to RDD"A").rdd.max()[0]

Each of the above gives the right answer, but in the absence of a Spark profiling tool I can't tell which is best.

Any ideas from either intuition or empiricism on which of the above methods is most efficient in terms of Spark runtime or resource usage, or whether there is a more direct method than the ones above?

6 Answers

+3 votes
by (13.2k points)

All the methods you have described are perfect for finding the largest value in a Spark dataframe column.

Methods 2 and 3 are almost the same in terms of  physical and logical plans. Method 4 can be slower than operating directly on a DataFrame. Method 1 is somewhat equivalent to 2 and 3.

There is another more efficient method, whose format is the same as method 3 . 


 Only difference from method 3 is that asDict() is missing.

If you wish to know about Hadoop Tutorial visit this Hadoop Certification.

by (19.9k points)
Very well explained. Thank you.
by (106k points)
Understood the concepts properly thanks
+2 votes
by (44.4k points)

A particular column's MAX value of a dataframe can be determined using this: 

max_value = df.agg({"any-column": "max"}).collect()[0][0]

by (29.3k points)
I would prefer this answer to be accepted.
Note: [0] [0] must be in the command to get the result.
by (62.9k points)
head() can be used instead if collect()[0]
+2 votes
by (108k points)


>row1 = df1.agg({"x": "max"}).collect()[0]

>print row1


>print row1["max(x)"]


The answer is almost the same as method3. but seems the "asDict()" in method3 can be removed.

by (47.2k points)
Max value for a particular column of a dataframe can be achieved by using -

your_max_value = df.agg({"your-column": "max"}).collect()[0][0]
+2 votes
by (32.1k points)

If you're looking to do this directly using Scala (using Spark 2.0.+), you do it like this:

scala> df.createOrReplaceTempView("TEMP_DF")
scala> val myMax = spark.sql("SELECT MAX(x) as maxval FROM TEMP_DF").
scala> print(myMax)
by (29.5k points)
I tried it like this .. this worked for me thankyou
0 votes
by (33.1k points)
edited by

In PySpark, you can do this simply by using this code:

max('ColumnName').rdd.flatMap(lambda x: x).collect())

For more information regarding the same, refer the following video tutorial:


0 votes
by (40.7k points)

Try using the code mentioned below:"A")).alias("MAX")).limit(1).collect()[0].MAX

For my data, I got the following benchmarks:"A")).alias("MAX")).limit(1).collect()[0].MAX

CPU times: user 2.31 ms, sys: 3.31 ms, total: 5.62 ms

Wall time: 3.7 s"A").rdd.max()[0]

CPU times: user 23.2 ms, sys: 13.9 ms, total: 37.1 ms

Wall time: 10.3 s

df.agg({"A": "max"}).collect()[0][0]

CPU times: user 0 ns, sys: 4.77 ms, total: 4.77 ms

Wall time: 3.75 s

Browse Categories