-
Notifications
You must be signed in to change notification settings - Fork 1
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
Merge pull request #9 from ryan-williams/ug
spark-util upgrade
- Loading branch information
Showing
12 changed files
with
95 additions
and
25 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -1 +1 @@ | ||
addSbtPlugin("org.hammerlab.sbt" % "base" % "4.6.1") | ||
addSbtPlugin("org.hammerlab.sbt" % "base" % "4.6.2") |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
26 changes: 21 additions & 5 deletions
26
spark/src/main/scala/org/hammerlab/cli/spark/PathApp.scala
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -1,21 +1,37 @@ | ||
package org.hammerlab.cli.spark | ||
|
||
import Registrar.noop | ||
import org.hammerlab.cli.base | ||
import org.hammerlab.cli.base.app.{ Args, ArgsOutPathApp, HasPrintLimit } | ||
import org.hammerlab.cli.base.close.Closeable | ||
import org.hammerlab.spark.confs | ||
|
||
/** | ||
* Generic Spark [[App]] | ||
*/ | ||
abstract class App[Opts]( | ||
_args: Args[Opts], | ||
reg: Registrar = noop | ||
)( | ||
implicit c: Closeable | ||
) | ||
extends base.app.App[Opts](_args) | ||
with HasSparkContext | ||
with confs.Kryo { | ||
reg.apply(this) | ||
} | ||
|
||
/** | ||
* [[HasSparkContext]] that takes an input path and prints some information to stdout or a path, with optional truncation of | ||
* such output. | ||
*/ | ||
abstract class PathApp[Opts](_args: Args[Opts], | ||
reg: Registrar = noop)( | ||
implicit c: Closeable | ||
implicit c: Closeable | ||
) | ||
extends ArgsOutPathApp[Opts](_args) | ||
with HasSparkContext | ||
with HasPrintLimit | ||
with confs.Kryo { | ||
extends ArgsOutPathApp[Opts](_args) | ||
with HasSparkContext | ||
with HasPrintLimit | ||
with confs.Kryo { | ||
reg.apply(this) | ||
} |
5 changes: 3 additions & 2 deletions
5
spark/src/test/scala/org/hammerlab/cli/spark/ConfigTest.scala
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
47 changes: 47 additions & 0 deletions
47
spark/src/test/scala/org/hammerlab/cli/spark/PathOptTest.scala
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,47 @@ | ||
package org.hammerlab.cli.spark | ||
|
||
import hammerlab.cli._ | ||
import hammerlab.indent.spaces | ||
import hammerlab.path._ | ||
import hammerlab.print._ | ||
import hammerlab.show._ | ||
import org.hammerlab.kryo | ||
|
||
class PathOptTest | ||
extends MainSuite(PathOptTest) { | ||
test("run") { | ||
val out = tmpPath() | ||
appContainer.main( | ||
"-n", "100", | ||
"-o", out | ||
) | ||
==( | ||
out.read, | ||
"5050\n" | ||
) | ||
} | ||
} | ||
|
||
/** | ||
* Test-[[Cmd]] similar to [[SumNumbersTest]], but taking its output [[Path]] as an option instead of as a positional | ||
* argument | ||
*/ | ||
object PathOptTest | ||
extends Cmd { | ||
case class Opts( | ||
@O("n") n: Int, | ||
@O("o") out: Path | ||
) | ||
|
||
case class Reg() extends kryo.spark.Registrar(classOf[Range]) | ||
|
||
val main = Main( | ||
new spark.App(_, Reg) { | ||
val out = opts.out | ||
out.mkdirs | ||
implicit val printer = Printer(out) | ||
val sum = sc.parallelize(1 to opts.n).reduce(_ + _) | ||
echo(sum) | ||
} | ||
) | ||
} |