Thursday, April 21, 2022

5 Problems Big Data Architectures Solve

 

The breadth of Big Data and Analytics technical architecture can seem intimidating. Despite its diversity, though, it solves a handful of problems. It makes trade-offs and surprise, just moves work around to meet its performance goals.

I’ll attempt to describe the key problems and the tradeoffs simply. These basics should help demystify architecture selection and troubleshooting for technology leaders new to the space.

Big Data is well, big. Too big for one machine to process quickly and store reliably. So, we need a lot of machines, working in parallel, and will have to coordinate their activity. How many? In the extreme, thousands.

Data may stream in so fast that it is hard to even ingest it without falling behind. With so much distributed data, finding and reading it quickly becomes a problem. Finally, reliable storage requires multiple copies and reliable coordination requires a way to handle coordinator failure.

In short, the problems are to:

1. Scale write, read, and processing performance

2. Parallelize work across nodes, i.e., machines.

3. Efficiently share data and coordinate activity across nodes.

4. Reliably process and store data despite one or more node failures.

5. Meet your business needs for data consistency and availability.

Problem 1: Scale write, read, and processing performance.

Writing

At scale, simple things like reading and writing become hard. The physical world is unforgiving, so learn your constraints. You can’t have it all. Ultra fast writes typically means slower reads and vice versa. When you can solve for both, accept that you’ll create a new problem. There are no free lunches at scale.

In a nutshell, store your data in ways that make it easy to do what you need to do with it. And address the consequences. Let’s look at some examples.

A traditional transactional database needs data consistency, so its relational model limits data redundancy. It needs reasonably fast reads and writes, so it defers some of its disc housekeeping, index and other storage optimization to off hours. Just don’t forget to schedule it!

Specifically, most Relational Databases record where they store their data on disc in real time using B+ Trees. Each file write causes multiple updates to the B+ Tree and takes time. To avoid more delay, it writes data somewhat haphazardly to disc blocks, slowing reads. The deferred work moves to clean up jobs. 

Source: Wikipedia

What if we need faster writes? Well, something has to give. Cassandra, for example, provides faster writes, strong read times, and lightweight transactions. The cost is more data redundancy and less consistency. You also have to agree to less flexible read options — sorts move from query parameter to design decision. But wow, very fast.

To do this, Cassandra, and others, e.g. Kafka, use a Log Structured Merge Tree. Writes are batched based on available RAM (memtables) and appended to a log (SSTables). Deletes and modifications are appended to the log rather than applied to the underlying data store. A compaction process handles data storage changes and optimizations, again in batches.

Source: Creative Coder

These are two common options. There are other optimization strategies that trade off various constraints.

Reading

Let’s stick with the familiar. A traditional relational database, Oracle, SQL Server, DB2, is row oriented, i.e., stores rows together in files on a drive. Rows offer a good fit for transactional use cases, lots of fine grained, surgical inserts, modifications, etc., and AVRO is a common file format. However, analytical use-cases frequently aggregate one column of data from thousands of rows and thus must read more than they need.

The heuristic is to store data the way you want to use it. Analytical use cases benefit from column oriented file structures, like Parquet. They read only what they need to and, given their columns, have homogenous data; so they compress to save storage space.

As always, there is a tradeoff: writes are slower. In fancy language, these structures are great for write once read many workloads. Column database examples, Redshift and HBASE.

Source: datacadamia.com

Yes, there are hybrids. Cassandra and BigTable can save a set of columns (a family), whose data is frequently used together.

You may also optimize for data access granularity. If you need to process large grained objects (files) one by one, then an object storage solution like AWSs S3 fits the bill. It uses a key-value pair where the value is the object, e.g. BLOB, Json file, etc. If you need fine grain access to groups of structured data elements, then consider a low-latency, high consistency, key-value/document database solution like AWSs DyanmoDB.

Note the storage types below. It fits this discussion well, but I’ll hand wave it and just provide a link so if you don’t already know, you can learn more about them. 

 

Source: Cloudain.com

Problem 2: Parallelize work across nodes, i.e., machines.

Since one machine, node, can’t do the work quickly enough, we must parallelize it. We need to abstract away the complexity of dealing with all the different nodes so the developer can focus on the business problem. Map Reduce is one such solution. How it abstracts away the complexity we’ll ignore and simply look at its workflow, which is basic.

Map Reduce does something simple. It collects some data and maps it into a structure, say sales for the week. And then reduces it into something useful, say the min, max, and average sale volumes for the week. The developer writes the logic for the job and the Map Reduce application finds the nodes with the data and runs the logic on each machine in parallel, and then stores it.

Source: Hadoop In Practice

While good for simple tasks, the Map Reduce framework causes a developer to create a lot of jobs and knit them together. Imagine quarter-end reporting as a long list of jobs. At Big Data scale, that is difficult to manage. As we’ll see, this is a solved problem.

Big Data Analytics require complex online analytical processing requires inner joins, outer joins, etc. Frameworks like Tez, HIVE abstract away that NoSQL/SQL complexity using algorithms expressed and coordinated via Directed Acyclic Graphs (DAGs). These graphs flow work from left to right and do not loop.

Source: Medium

DAGs are expressive and easy to use and have become a de facto developer standard for expressing simple and complex workflows, i.e., they are not just powering low level frameworks. Examples include Spark, Beam, Airflow, Kubeflow, and ML pipelines…

ML uses interesting graphs. Feed-forward Neural Networks traverse a graph in one direction, have layers, and evaluate loss at the points where layer nodes meet. Convolutional Neural Networks, and others, use bi-directional traversal, to enable Gradient Descent and Backpropagation.

Source: Brilliant.org 

Problem 3: Efficiently share data and coordinate activity across many nodes.


Many machines working in parallel mean many logical and physical boundaries to cross, process, user space, system space, network, and storage mediums, etc. Given speed is a goal, boundaries are a toll, they take time to cross.


Process as much data locally as is practical, I.e. minimize the amount of data passed across the network, processes, or storage boundaries on a node to reduce latency. As mentioned earlier, Map Reduce handles some of this heavy lifting. However, remember it passes data across jobs by persisting it.


As an example, in the Hadoop ecosystem, YARN applications may keep nodes (VMs) active across jobs to keep data in memory for a subsequent workflow step and to avoid VM startup times. The tradeoff is inactive VMs consume memory and may starve others. And there are other examples, Spark uses Resilient Distributed Datasets (RDDs) as a form of shared memory.

Source: tutorialspoint.com

Stating the obvious in distributed systems, avoid workflows that allow hosts to block each other’s tasks, and prefer asynchronous vs synchronous cross node communication. Onto coordination.

Coordination

Coordination happens at multiple levels. We'll keep it simple. Let’s define “coordinate” to mean managing work across nodes. The next section will address failures.

Continuing to use the Hadoop example, let’s look at the “abstracted complexity” that allows developers to not worry about work distribution. There are three layers:

1. YARN (Yet Another Resource Manager) - yep, manages applications, handles some scheduling.

2. MapReduce - organizes where mapper and reducer work will happen (runs in YARN as of Hadoop 2).

3. HDFS - Hadoop Distributed File System - determines where data goes where it is.

Source: Hadoop In Practice

The above diagram highlights a common coordination model, master/slave, or leader/follower, if you prefer. Here, the leader is the authoritative source of what data is “correct” and which node should do what and when. Leader/follower relies on an primary/backup data replication model. A set of changes comes in, let’s say to add a value, and then update it three times. The leader will apply those changes and send the final value to its followers. This can create a single point of failure, so often a backup master who is aware of all the changes. As long as both of them don’t fail at once, all is well. Over two masters? Then there are some overlaps with peer-to-peer. 

In centralized distributed computing, leader/follower and peer-to-peer are the main options.

In peer-to-peer, any node can be a leader or a follower; all nodes have the same capabilities and their function, leader or follower, is context dependent. Here, the data replication model is state machine based. While there is a leader, it shares all incoming changes with the other nodes in its cluster so that anyone of them can become the leader if the current leader fails. 

Source: The Log: What every software engineer should know

Source: Tutorialspoint.com

Cassandra uses a peer-to-peer model in a ring topology for replication. Its distribution model is blissfully simple, racks and data-centers. Automated replication across data centers is a popular feature. To be nerdy, very cool.

Source: Intstacluster.com

Your work load should drive your choice. Leader follower is simpler, is easier when strict consistency is required, and with back-up masters is reliable. However, with enough inbound requests, it may bottleneck.

Peer to peer is more complex, is chatty (i.e. requires more internode communication), and prefers eventual consistency. However, this complexity provides dramatic horizontal scalability. For example, a Cassndra deployment ran12,000 commodity servers across many data centers/regions. When deployed properly, you won’t get diminishing returns, it scales linearly.

Problem 4: Reliably process and store data despite one or more node failures.

At scale we coordinate 100s to 1000s of nodes and serve 1,000,000s of customer. Some of them, hopefully only the nodes, will fail, and sometimes many nodes will fail together. So peers, followers, and leaders all need to share or replicate their data to at least one partner, and frequently 3 to 5 partners. This data might be a customer order or the identity of the current leader on the network when one has failed.

Regardless of the distribution style, node failure in an asynchronous system introduces a hard constraint - it is not possible to be 100% available and provide 100% data consistency. Yes, those are strong words. Read the FLP proof here

Given three or more copies of every piece of data, that data values change, and that we can’t instantaneously update all nodes, which nodes have the correct value? And what happens when a node fails? Did only it have the correct value? How would you know? Again, we’ll keep it simple here by discussing only a common conceptual approach, the Paxos algorithm.

The problem is achieving consensus about the correct value of a piece of data at a particular time across multiple nodes when one node can fail. Thankfully, the terminology is straightforward: the nodes must come to a consensus on the correct value, say the value of items in a shopping cart, or perhaps the identity of the current leader.

The lowly log, append only, is often at the center of achieving consensus during replication in this asynchronous world. It has a lot to do with the order in which things happen. But we’ll leave that for another day.

The Paxos algorithm in the abstract is simple. Each participant is a node. A node proposes that the other nodes promise to consider a new value for a shared data element. If a quorum of other nodes promise to consider accepting a new value, then the proposer sends out that value. If a quorum of nodes accepts the new value, then the proposer tells all the nodes to commit the value. And how does each node decide? Well, Paxos is a paper unto itself, so read about it here

Cassandra, S3s indexing solution, and many others use Paxos. Other solutions, like Zookeepers atomic messaging, used by Kafka, Hadoop, etc. exist as well. 

Problem 5: Meet your business needs for data consistency and availability.

The CAP Theorem is important to understand when talking about distributed systems. It is often misunderstood. CAP stands for Consistency, Availability, and Partitioning. Think of partitioning as a failure, a failure of either a node or a network connection. 

The point of the Theorem is that, since we can’t avoid node failure (partitioning) in a distributed system, we must optimize for either consistency or availability. The key word is “optimize.” An AP system can be consistent through failures, but it may take longer than your customer can accept.

CP System - prioritizes consistency over availability/speed, i.e., I’m willing to wait because I need the most current data

AP System - prioritizes availability/speed over consistency, i.e., I’m not willing to wait. Give me what you have now. I’ll deal with it.

If you must guarantee consistency and Availability AP, then don’t use a distributed system as it may fail, but it can remain consistent.

Source: Medium


Thursday, April 7, 2022

Put simply, how does Supervised Learning work?


In Supervised Learning, a program ingests training data as sets of observations, called features, identified using labels. For example, many emails, each labeled Spam or Not Spam.

The program expresses the features mathematically and sends them iteratively into a function. It varies the functions’ governing parameters until getting the desired output. Emails correctly identified as spam or not spam. The program is a machine learning model and the iterative process is called training.

Programs, like people, are not perfect learners. Data scientists and engineers evaluate a model by identifying loss, the number of mistakes made, and variance, how well the model performed across different sets of training data.

Too little loss may cause overfitting where model results are erratic, i.e., have high variance. Too much loss yields more consistent results across data sets, but not accurate enough to be useful.

Data scientists send their programs a lot of data and use a lot of computational resources and time to enable successful machine learning. Sometimes their math enables an accurate model that runs in milliseconds on many data sets, for example, voice recognition on your phone. Sometimes their mathematical approach just fails. Very much a case of “if at first you don’t succeed, try, try again.”

The math and the technology involved are non-trivial. Today’s data scientists stand on the shoulders of centuries of pragmatic mathematicians. Moore’s law and distributed big data together enable the computing scale required.


Friday, March 11, 2022

Shape Up, Agile Method Summary and Commentary

 


ShapeUp! Shaping, Betting, and Building

The summary starts with the second step in the method, the story flows better this way,

Betting on a six-week release

Key concepts:

  • Product and engineering commit to share and mitigate delivery risk (Fix time/effort and vary scope.)

  • Cross-functional, autonomous teams align to independent technology components, e.g., services, and have end-to-end accountability for feature design, delivery, and production operations.

Phase Overview:

A 6-week release effort is a “bet” matching an “appetite”

An “appetite” is not an estimate. Estimates start with a design and end with a number. Appetites start with a number and end with a design. The appetite is a creative constraint on the design process - to keep it in check, balanced to value.

  1. Language connotes a business risk worth taking - clear customer value at a reasonable price

  2. Needs mutual commitment from tech and product to vary scope and approach to win the bet.

  3. Caps the downside: short enough time to limit the damage if it will cost more than its worth

  4. Provides pressure: long enough time to get something meaningful done and short enough to feel the date pressure

Product shapes the next idea as they support the build process of the current idea.

Shaping

Key concepts:

  • Shape an idea for a customer outcome to design and build in one release.

  • Set boundaries, identify risks, and layout a high-level model, not a design, to be elaborated during the build

  • Pitch the idea to place a bet, i.e., to be chosen to attempt delivery

Phase Overview:
  1. Start with a raw idea - What problem does it solve? And what outcome gives it customer value? How will we verify they get it?

  2. Shape the idea to fit an appetite, apply design thinking, i.e., it may need to be decomposed, list the constraints.

  3. Set boundaries - how much is enough?

  4. Rough out the elements of the idea at a high level, low-fi, but clear on the outcome. Breadth, not depth — explore options. Leave room for designers, e.g., not a UI spec.

  5. Address risks, and rabbit holes by looking for unintended consequences, unanswered questions, etc. Specify the tricky details.

  6. Get technical review and determine what is out of bounds.

  7. Write the pitch:

    1. The problem to solve, along with the expected customer outcome and verifier.

    2. Our Appetite - how much time is it worth and what constraints does that imply?

    3. The core elements of the solution - not the “answer.”

    4. The Rabbit holes to avoid and risks to deal with.

    5. No-gos - what should the team exclude, things we are choosing not to cover to fit the appetite or make the problem workable.

Building

Key concepts:

  • Apply design thinking to balance feature design, technical risk, and time to market.

  • Organize work in the team by application structures (“scopes”) not people. Scopes are independently buildable and testable and may depend on each other.

  • Do the hardest/riskiest thing early

Phase Overview:
  1. Product assigns projects, not tasks, and done = deployed

  2. Hand delivery over to the team to build a feature that gets the outcome given the technology and the time available.

  3. Discover and map the scopes, the independently testable and buildable, end-to-end slices that together make up the feature. Use these scopes, e.g., edit, save, send, to show progress.

  4. Get one piece done, a small end-to-end slice to gain momentum within a few days

  5. Start in the middle, with the most novel, risky element. If time runs short, simplify or remove nice to haves or should-haves.

  6. Substance before style - build and verify basic interactions work before focusing on UI styling

  7. Unexpected tasks and opportunities will appear as you go, so know when to stop.

    1. Compare completed work to a baseline, e.g. the customer experience now, not a future ideal.

    2. Use the mutual commitment of 6 weeks to an all-or-nothing release as a circuit breaker to limit the scope.

Commentary

Top takeaways:

  • Scopes with automated tests speed development and enable the long run product and organizational flexibility that is central to Amazon, Google, Spotify, etc.
  • The “circuit breaker” motivates frequent releases and shared accountability between engineering and product, but may trade-off completeness in the near term.
  • ShapeUp is lightweight and has obvious limits. Say you’ve got 48 people across 8 teams and two years of budget to scale up your software. How do you define and coordinate all of that work? Carefully, I assume.
  • Shape Up heroically assumes autonomous, cross-functional teams. Which I support wholeheartedly, but you may not have.
The concepts are excellent and apply outside this lightweight method. I added the italicized content as it felt implied. I’m guessing most who use it would also sprinkle in some scrum. This approach begs for XP practices. Scopes are natural outcomes of TDD and BDD.

The summary leaves a lot out, e.g., large project how-to guidance. Some of it seemed silly, e.g., visualizing status as scopes rolling up and down hills. But, the content is free; you find it here, and I’m not complaining. The scopes concept is central to software development. In fact, outside of the hill thing, there is little to dislike.

ShapeUp is lightweight and has obvious limits. Say you’ve got 48 people across 8 teams and two years of budget to scale up your software. How do you define and coordinate all of that work? Carefully, I assume. It heroically assumes autonomous, cross-functional teams. Which I support wholeheartedly, but you may not have.

A small consultancy doing web development projects for clients created the method, and it fits that like a glove. It can be a great fit for small tech startups until they scale past two or three teams. After that, keep the concepts and solve the next set of problems. Lots of good toolsets out there, LeSS, Scrum, etc.


Friday, March 4, 2022

My most popular content...

 

Prior to 2022, google analytics shows my most popular content is the book summaries. So, here's a list with direct links. 

  1. The DevOps Handbook - a four-part book summary
  2. Lean Start-Up - a three-part book summary
  3. Kanban - a seven-part book summary
If you like this type of content then you are sure to enjoy these as well:
  1. The Nature of Software - Mary Poppendieck (30-minute read)
  2. Agile Fluency - The Agile Fluency Project
  3. Large Scale Scrum - the LeSS framework by Craig Larman & Bas Vodde
If you are looking for an afternoon read, Doc Norton's book Escape Velocity is excellent. A short book but packed with good software delivery practice wisdom. 

Have fun!

Saturday, January 21, 2017

Has Agile Lost It's Way?



When I got started with Agile in 2007 it was new. Most of us had been delivering software for a while and while we loved developing solutions we hated the often tyrannical circumstances we endured doing it. We knew what worked: lots of unit testing, iterative/incremental development, lots of collaboration, transparency, etc. But until Agile (for me, XP, and Scrum), we didn't know how to pull it all together. What we found, before Daniel Pink gave simple ways to say it, was that when we had the opportunity to 'do it right,' we were on a path to mastery, autonomy, and purpose in our daily work.

Agilists associate the birth of "Agile" with The Snowbird Conference in 2001 and The Agile Manifesto. If we focus on the unmet needs of businesses and governments relying on software then it goes farther back still.  Let's just say Agile's about 16, a teenager. It shows.

Agile has become a big business promising homogenized faster, better, and cheaper delivery to the masses. I would guess that 70% or more of teams who call themselves Agile aren't. At least by my definition but, I'm a tough grader who expects the use of XP engineering practices.

In 2012 Thoughtworks' published an Agile Fluency Model to describe levels of Agility from One Star Teams focused on delivering business value to Four Star Teams optimizing entire eco-systems. Their data suggested that 15% of teams were doing something, but it wasn't Agile. 45% of teams, while focusing on business value, weren't meaningfully improving software quality. As a result, due to Scrum or Scrum-like processes, these teams are only marginally more productive.


As a result of this homogenization, the tyrannical circumstances Agile was to dispel are back for too many. As the market grew, "consultants" with no more than a scrum class and a half-read, frequently derivative book popped up to feed on the Agile sales frenzy, the essence started getting lost.

Turns out that greater transparency can also lead to micro-management and ever more unrealistic expectations on a team. You get worse, not better code. You get less not more shared understanding. Even when you've had some success by focussing solely on visible business value, you may hear, "Wow, we've gotten faster since leaving waterfall. Bet you can go twice as fast! Indeed, you must! Find more process improvements!"

If this is sounding familiar to you I'm sorry. Here's what I suggest you might reference to find the essence:
  1. Read the original sources:
  2. If you are a developer or architect then also check out:
  3. If you are an Architect add these to the above:
    • Read Clean Code
    • Watch the Architecture, Tech Debt, and TDD videos at Clean Coders
    • If you haven't written a line of code in 10 years, and think this is all craziness then spend the next month coding, full-time in your enterprise applications then read and watch the videos again.
Once you are done, happy New Years' 2006! Please don't stop now. Next up:
  1. RSA ANIMATE: Drive: The surprising truth about what motivates us
  2. Kanban: Successful Evolutionary Change for Your Technology Business
  3. The Lean Startup: How Today's Entrepreneurs Use Continuous Innovation to Create Radically Successful Businesses
  4. Large-Scale Scrum: More with LeSS
  5. The DevOps Handbook
  6. BADASS: Making Users Awesome
Ok, happy New Years' 2017! We're not done. So much has happened since Agile was born. Not all of it captured above. Keep in mind that you have to know the rules before you can successfully break them. You should also check into Modern Agile



Friday, January 13, 2017

Finally Finished the Kanban Book Summary!

Yes! Finally finished the seven-part summary of "Kanban: Successful Evolutionary Change for Your Technology Business by David Andersen. Part one starts here.



Of all the summaries so far, this is one of my favorites. The summaries are written as study guides or review materials.

The Kanban book fills in the blanks regarding why Scrum and XP work the way they do and how to improve their processes in a way that minimizes resistance to change and maximizes value delivered.

Kanban is not a methodology for software delivery. It is a change management system for improving the implementation of any delivery framework or process. It is reasonable to say that using the system to improve a waterfall process would lead to Scrum/XP/Lean Start Up like processes. This may explain why some think Kanban is a delivery methodology - applying its principles will suggest process improvements that likely converge towards the processes that are taken to define today's "Agile."

But that too is perhaps misleading. Agile is a mindset growing from of a set of principles and values we apply to improving our ability to deliver value through software delivery. Is is not a process or a framework - despite industry marketing machine claims to the contrary. It uses processes and frameworks and it is our job to improve them. If you agree then you'll enjoy getting involved with Modern Agile.

David Andersen's application of Lean principles to software delivery and their synthesis with the Agile mindset makes Kanban a valuable tool. Here are the summaries: