Friday, June 12, 2020

Spark, Parquet and Schemas


Recently, I've been playing with Spark and Parquet. Here are a few tidbits I thought were interesting.

Nulls

I was surprised to see that attaching a Schema to a DataFrame did not enforce its constraints ("Schema object passed to createDataFrame has to match the data, not the other way around" - SO). Instead, Spark does a "best efforts" attempt to have data read from HDFS to comply with the schema, but it does not throw exceptions if it can't be done. For instance, columns in the schema that don't appear in the file are filled with nulls even if the schema says they cannot be null.
"When you define a schema where all columns are declared to not have null values , Spark will not enforce that and will happily let null values into that column. The nullable signal is simply to help Spark SQL optimize for handling that column. If you have null values in columns that should not have null values, you can get an incorrect result or see strange exceptions that can be hard to debug."
Spark: The Definitive Guide: Big Data Processing Made Simple (p102)
The solution was to deserialize everything into the appropriate case class for the Dataset. Since I had to do bespoke validations on the objects, this wasn't too great an issue.

Note that "when writing Parquet files, all columns are automatically converted to be nullable for compatibility reasons" (from the Spark docs). The compatibility being the Parquet spec (see below) which we have to map to when writing.


Floating points

Reading a parquet file that had decimal(38, 10) for a floating point type, taking its schema with DataFrame.schema and then using that self same schema to load the Parquet again with:

spark.read.schema(theSchema).parquet("...").as[MyDomainObject].show()

unfortunately gives me:

org.apache.spark.sql.AnalysisException: Cannot up cast `REDACTED` from decimal(38,10) to decimal(38,18) as it may truncate
The type path of the target object is:
- field (class: "scala.math.BigDecimal", name: "REDACTED")
- root class: "MyDomainObject"
You can either add an explicit cast to the input data or choose a higher precision type of the field in the target object;

where MyDomainObject has a field that is simply a scala.math.BigDecimal.

The problem seems to be due to the underlying store from whence the data was plucked (as this SO question shows a similar problem with MySQL). 

Modifying the schema by hand to use only decimal(38, 18) solved the problem. This inconvenience has its own JIRA that's unfortuately marked as "Won't Fix".

Since "the types supported by the file format are intended to be as minimal as possible" (from the Parquet docs), the floating point types are your standard IEEE formats. It appears that if you want BigDecimals, the data will be saved in a Java representation by using Parquet's native BYTE_ARRAY types since Parquet is language agnostic. Spark assumes as it's defined as a decimal(38, 18) in SYSTEM_DEFAULT here.

Since this is the default, if you save your data in anything less, you must read the file with a schema if you're hoping to have a Dataset.as[MyDomainObject]. You might need to cast the original values as this StackOverflow answer suggests or else the decimal point might not be were you were expecting it to be.


Timestamps

As already mentioned, the Parquet types are restricted to a IEEE compliant types and a few primitive others. Therefore, all the richness of the Java ecosystem is lost and there needs to be a mapping between these primitives and Java classes. It appears that Dates and Timestamps are converted from the INT64 type as you'll see casting errors if your types misalign. 


Sharing schemas

You can load and save the schema as JSON if you like (as this SO answer shows). You can also add metadata to the schema with a simple:

import org.apache.spark.sql.types._
val metaBuilder = new MetadataBuilder()
val sillyFields = schemaFromJsonFile.map(_.copy(metadata = metaBuilder.putString("You", "suck!").build))
val sillySchema = StructType(sillyFields.toArray)
spark.read.schema(sillyStructType).parquet("...").show()

and all is good with the World.

The advantage of this is that you can pass annotated schemas between teams in a (fairly) human-readable format like JSON.


Conclusion

Spark schemas and Parquet do not align one-to-one since Spark schemas could apply to any file format (or even just data in memory). Some elements of the schema are lost when written to Parquet (eg, everything in Parquet is nullable irrespective of the Spark schema) and some types in Scala/Java are not represented in Parquet.

Saturday, June 6, 2020

Cracked Pipe


Pipes can get blocked as an old post of mine illustrates. So, how do we ensure that we can read and write to and from pipes in the effectful world?

Well, my first (rather poor) attempt was to use Java's PipedInputStream and PipedOutputStream and have ZIO handle the multithreading aspects of it. Although the test passes it is far from exhaustive.

One of the problems is that unbeknown to me that Java's piped IO classes are very thread sensitive. Consequently, my ZIO test was failing with:

    Fiber failed.
    A checked error was not handled.
    java.io.IOException: Read end dead
        at java.io.PipedInputStream.checkStateForReceive(PipedInputStream.java:262)
        at java.io.PipedInputStream.awaitSpace(PipedInputStream.java:268)
        at java.io.PipedInputStream.receive(PipedInputStream.java:231)
        at java.io.PipedOutputStream.write(PipedOutputStream.java:149)
        at java.io.OutputStream.write(OutputStream.java:75)
        at uk.co.odinconsultants.fp.zio.streams.LargePipeMain$.writing$1(LargePipeMain.scala:42)
        at uk.co.odinconsultants.fp.zio.streams.LargePipeMain$.$anonfun$piping$3(LargePipeMain.scala:58)
        at zio.internal.FiberContext.evaluateNow(FiberContext.scala:458)
        at zio.internal.FiberContext.$anonfun$evaluateLater$1(FiberContext.scala:687)
        at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
        at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
        at java.lang.Thread.run(Thread.java:748)

This happens because the JDK code has something in it that looks like this:

...
        } else if (readSide != null && !readSide.isAlive()) {
            throw new IOException("Read end dead");
        }
...

and to be fair, the Java API documentation says:

"A pipe is said to be broken if a thread that was providing data bytes to the connected piped output stream is no longer alive."

But if your effect system uses fibres instead of threads then it is free to stop and start threads as it pleases. The bottom line is that wrapping side effecting code in a ZIO or a Cats IO might not be the first step in making your code more FP.

Saturday, May 23, 2020

When Penguins Attack


It's that moment that everybody who auto updates their laptops fears: it doesn't work. My problem was that after logging in, the OS paused briefly and then sent me back to the login screen. Hmm.

In this situation, get a CLI prompt with CTRL-ALT-F3. Since we managed to get a login prompt and then failed, it's probably the video drivers. I checked and saw a message similar to this (segmentation fault) in /var/log/Xorg.0.log

Foolishly, I chose to rollback a recent kernel update by updating /etc/default/grub. This is a matter of setting GRUB_DEFAULT to a number indicating how many steps back you want (0 obviously uses the most recent). Don't forget to run sudo update-grub when you're done.

This was foolish as I had rolled back to a kernel version that before a recent BIOS update. Now, things were much worse as my laptop could not get past the BIOS screen. Great. Even trying holding down shift at boot time to allow me change kernels didn't work.

If you've bought a Dell (I have a Precision 7540, an ugly but decent machine), you'll be given a System Service Tag on the the invoice they've sent you. Give Dell this ID here and download a recovery ISO. I tried phoning Dell Support but they never got back to me.

Create a bootable USB using Startup Disk Creator. Naturally, you'll need another Ubuntu machine to do this.

As  it happens, this startup disk just tries to install another version of Ubuntu which blew a bit of disk space but there was no reason to do this. I just dropped to a CLI with CTRL-ALT-F3 and saw what hard disks there were with fdisk -l, identified the partition that was clearly mine (at over 800gb) and mounted it with:

mount /dev/MY_DISK  /mnt/old

(this SO answer proved invaluable). Now I can read my old disk. Good.

Now, I can rollback the change I made to my grub file and sudo update-grub again. But first, I needed to chroot that mounted disk to pretend it's running as my root file system. Don't forget to also mount other important mount points to make it look genuine. This SO answer will help.

Rebooting got me back to my original problem of not being able to run XServer. The solution was: keep the kernel and rollback the graphics driver updates with a third SO answer here. For me, this was sudo apt-get purgeing all things NVIDIA and running sudo apt-get install nvidia-driver-430 (the driver version may be different for you).

Rebooting once more got my system back but c'mon you guys at Dell and Ubuntu - don't break my system. That's a few hours I won't get back.

Friday, May 8, 2020

Compiler settings for bug hunting


Here's a subtle bug that tripped me up this week.

Take this innocuous ZIO test code:

  override def spec: ZSpec[TestEnvironment, Any] = suite("Look closely....")(
    testM("Why is an obviously wrong assertion passing?") {
      ZIO("Actual").map { x =>
        assert(x)(equalTo("Actual")) &&
          assert(x.toUpperCase)(equalTo("ACTUAL")) &&
          assert(x)(equalTo("Obviously wrong")) &&
          assert(x.length)(equalTo(6))
          assert(x.toLowerCase)(equalTo("actual"))
      }
    }
  )

The test should fail as x is clearly not "Obviously wrong". I wasted an hour wondering whether I really understood the ZIO effect system at all. 

Then I saw the missing  && d'oh. 

What can you do to mitigate such errors? This is one kind of bug that wouldn't happen (or would at least be obvious) in Python where whitespace is significant. Wouldn't it be good if the compiler caught this?

Well, you can get the compiler to emit very loud warnings by using the warn-value-discard flag. You can make the compilation fail with -Werror or -Xfatal-warnings (see here). And you can fine tune the warnings with Scala 2.13.2 (see here).


Monday, April 20, 2020

Graph databases and Data Collisions


Back in 2016, we were comparing Neo4J and OrientDB graph databases. I was tasked with exploring OrientDB and engaged the CEO, Luca Garulli, to help me bulk load a huge amount of data (our chat can be found here).

In both cases, there is a problem with inserting edges in that each edge insertion can update two vertices. Let's assume there are no edges connecting the same vertex to itself (there really isn't in our data), so each insertion updates exactly two vertices. Well, if two edges are inserted at the same time there is the possibility of a collision.  This is a version of the famous Birthday Problem.

We were advised by the OrientDB guys to batch our edge insertions. But what's a good batch size? A high number would mean greater throughput but also a higher chance of collision that would ultimately force us to retry.

Let's say there are N vertices that have already been written to the DB and our batch size for writing edges is b. What are the chances of a collision for a given batch of edges?

Firstly, we need to find the expected number of collisions within a batch of size b. These are OK as they won't force us to retry since they are part of the same transaction.

Here, we note a nice trick outlined here.
The expectation of a sum is the sum of the expectations. This theorem holds "always." The random variables you are summing need not be independent.
So, the expected number of collisions, E(X) is

E(X) = E(collisions for first vertex) + E(collisions for second vertex) + ...

The expected value of anything is the sum of all possible values multiplied by their probability. So, the expected number of collisions for each vertex is

E(collision of vertex x) = 0 * p(no collision) +  1 * p(collision)  = p(collision)

Note that the expected value need not be an integer nor indeed sensible. The expected number when throwing a 6-sided die is 3.5.

Given n vertices, this number becomes:

E(X) = bp

where the probability is given in the link above, that is

p(collision) = 1 - (1 - (1/N))b-1

For anything reasonable, this is pretty small. For instance, if there are 1 million vertices and we choose a batch size of 1000, the expected number of collisions within a batch is just less than 1.


Monday, April 13, 2020

Applicative and Effectful Types


Typical use case: we have several legacy java.io.OutputStreams to close after we have done some work. We want to try to close all of them even if some might fail to close. Can we encapsulate these individual units of work in monads?
Recall a key distinction between the type classes Applicative and Monad - Applicative captures the idea of independent computations, whereas Monad captures that of dependent computations. Put differently Applicatives cannot branch based on the value of an existing/prior computation. Therefore when using Applicatives, we must hand in all our data in one go. [Cats documentation]
So, not only is monad inappropriate, it's impossible to perform this use case with monads at all.

Ok, let's look more closely at these Applicatives:
Rob Norris @tpolecat There's already a common misconception that applicative composition implies order-independence, which isn't necessarily (or even commonly) true. I feel like we're still kind of struggling for a way to talk about effects. I am, anyway.
Fabio Labella @SystemFw Because most education resources state that Applicative does represent non ordered computation whereas it can represents non ordered computation.
Rob Norris @tpolecat Even Validated is order-dependent
[From Gitter]

My reading of this is that calling Apply's *> will "Compose two actions, discarding any value produced by the first." (Cats Apply Scaladoc). The order is important in that all but the last result is thrown away.

To demonstrate this, I wrote some code here [GitHub] that uses Validated from Cats which is an Applicative as its code [GitHub] proves here.

My code uses a filthy var but as even the esteemed Fabio Labella "you need a Ref, but not a Ref of aWriter, just a Ref with IO. You can also use a var, whether you use one or another often depends on how complex the test is (var works well for simple cases with no concurrency)". So, with that caveat, let's proceed.

      var i = 0
      def failure(): Validated[String, String] = {
        println("failure")
        i = i + 1
        Invalid("failure message")
      }
      def success(): Validated[String, String] = {
        println("valid")
        i = i + 1
        Valid("success message")
      }

      success() *> failure() *> success()
      i shouldBe 3

This passes. Now, let's look at a monadic version in the same test class:

    type     MonadType = Either[String, Int]
    val aye: MonadType = Right(1)
    val nay: MonadType = Left("nope")

    aye *> nay *> aye shouldBe nay

You can demonstrate if you like that the last aye is not actually called but this is left as an exercise.


ZIO

ZIO is another Effects system. It is totally independent of Cats (unlike Monix, yet another Effects system, which depends on Cats).
"A quick summary: efforts to use Pure Functional Programming in Scala began with Scalaz lib which was then continued by Cats. These libraries are great, but in essence they try to replicate the whole Haskell experience in Scala. This comes at a cost mainly because Haskell is a non-strict, lazy-by-default language while the JVM is both eager and strict, so in order to make some things work some techniques were used that affect your performance and Scala's ability to infer types.
ZIO strives to give you a PFP experience but taking adavantage of Scala's paradigms (such as variance) to help both with performance and type inference.
it also tries to be more friendly for newcomers" ToxicaFunk (12 April 2020) on Discourse
Now, although all monads are applicatives, not all applicatives are monads (cats.data.Validated being an example of such a type). "We sometimes use the terms monadic effects or applicative effects to mean types with an associated Monad or Applicative instance." [Functional Programming in Scala]

In ZIO, there is "one monad to rule them all", also called ZIO. So, how do you get Applicative behaviour? None other than ZIO's creator himself answered me on Discourse:
jdegoes 12 April 2020 at 2:10 PM
@PhillHenry x1.ignore &> x2.ignore &> x3. This will execute x1, x2, and x3 in parallel, using zipRightPar (that's the &> operator, zips two effects together in parallel, returning whatever is produced on the right), and ignore the result of x1 and x2 so their failures don't influence the result of the computation.
@PhillHenry In ZIO, even parallel zip or collect or foreach operations will "kill" the other running effects if one of them fails. Because that's often what you want. To get the other behavior, just use .ignore in the right places to ignore the failures you don't care about.
My equivalent ZIO code can be found here on GitHub. Looking at the thread names and stack traces, these operations do appear to execute on different thread/fibres. This is not the case on the Cats implementation.

Wednesday, April 8, 2020

Mathematical Proofs for Distributed Systems


Software is becoming more about proofs. With well written FP code, you can prove at compile time that your software cannot enter certain pathological states.

Well, the world of architecture is moving in that direction, too.


Introducing TLA+

TLA+ is an attempt to help architects and developers to write bullet-proof systems. It's particularly useful in the domain of distributed computing.
"Distributed algorithms and protocols are hard to get right, especially, when they have to tolerate faults. One of the reasons is that distributed algorithms vary in their assumptions about thedistributed system: the communication medium, system synchrony, possible faults, etc. ... [TLA] offers a rich syntaxfor sets, functions, tuples, records, and sequences on top of first-order logic"
TLA+ Model Checking Made Symbolic

It's already being used in architecting Kafka.


Propositional Logic

In propositional logic, "there are only two values, TRUE and FALSE. To learn how to compute with these values, all you need to know are the ... definitions of the five Boolean operators"[1]: ∧, ∨, ¬, ≡ and ⇒.

The first four are AND, OR, NOT and EQUALS, so nothing special there. The last one (the implies operator) is tricksy. Its truth table looks like this:

FGF⇒G
TRUETRUETRUE
FALSETRUETRUE
FALSEFALSETRUE
TRUEFALSEFALSE

To help understand it, let's look at this example formula:

(n > 3) ⇒ (n > 1)

The first three rows of the truth table simply describe n>3,  1<n≤3 and n≤1.

The last row is when n>3 but not n>1 which is clearly wrong, hence FALSE.

Now, say you were asked to prove that:

(F⇒G) ≡ (¬F∨G)

is a tautology. You could write out all the truth tables which would be laborious. "However, computers are better at doing this sort of calculation."[1]


Very nice, but what's the application?

You can find a very nice model of some Kafka functionality here on GitHub. You can open this in the IDE called TLA+ Toolbox. It looks like this:


The language is pretty straight forward. For instance, pre-pending formulas with operators like /\ is syntactic sugar for removing wide parenthesis and EXTENDS seems to be like import in Java and Scala. You can find some blogs to learn more about it here [LearnTLA], here [Anton Sookocheff's blog] and here [Jack Van Lightly's blog].

The Toolbox can turn this into a more readable PDF (more readable for mathematicians anyway):


You can then "run" your model by giving it initial parameters. The Toolbox will then explore its state space. Interestingly, Hillel Wayne in this video descibes the path through the states as a description of the behavious of a system.

Another video can be found here and the transcript to Lamport's own presentation can be found here.

Conclusion

TLA+ looks like it could turn the IT architecture industry into something more than "it feels kinda right" that we see at the moment. One only hopes that its adoption in projects like Kafka will encourage other people to learn it. Me - I've only taken the first few steps. 

[1] Specifying Systems, Dr Leslie Lamport