HNHacker News
TopNewBestAskShowJobs

robertnishihara

464 karma · joined October 16, 2017

submissionscomments
robertnishihara··on Amazon's exabyte-scale migration from Apache Spark to Ray on EC2
To clarify, what I mean is that working with "exabytes" is atypical. Most use cases are at a slightly smaller scale :)

Data processing workloads are quite common on Ray, especially with unstructured data.

Also, I work on Ray, which is the underlying framework used here, but all the work in the post was done by the Amazon team.

robertnishihara··on Amazon's exabyte-scale migration from Apache Spark to Ray on EC2
Other folks have built data processing libraries on top of Ray: Modin and Daft come to mind.

But I'm not aware of anything exactly like what you're referring to!

robertnishihara··on Amazon's exabyte-scale migration from Apache Spark to Ray on EC2
Multi-threaded libraries (e.g., numpy and PyTorch on CPUs come to mind) are well supported. In scenarios where many processes are each running heavily multi-threaded computations, it can help to pin specific processes to specific cores (e.g., using tools like psutil) to avoid contention.

The scenario where a Ray task forks is probably not very well supported. You can certainly start a subprocess from within a Ray task, but I think forking could easily cause issues.

You can definitely use Ray + Jax, but you probably need to avoid forking a process within a Ray worker.

robertnishihara··on Amazon's exabyte-scale migration from Apache Spark to Ray on EC2
I'm glad you find it exciting!

Our intention from the start was for Ray to be general purpose. And the core Ray APIs are quite general (basically just scheduling a Python function somewhere in a cluster or instantiating a Python class as a process somewhere in the cluster).

We had AI use cases in mind from the start, since we were grad students in AI. But the generality has really been important since AI workloads encompass a huge variety of computational patterns (allreduce style communication patterns on GPUs for training, embarrassingly parallel data processing workloads on spot instances, and so on).

robertnishihara··on Amazon's exabyte-scale migration from Apache Spark to Ray on EC2
Yeah, mmap, I think this is the relevant line [1].

Fun fact, very early on, we used to create one mmapped file per serialized object, but that very quickly broke down.

Then we switched to mmapping one large file at the start and storing all of the serialized objects in that file. But then as objects get allocated and deallocated, you need to manage the memory inside of that mmapped file, and we just repurposed a malloc implementation to handle that.

[1] https://github.com/ray-project/ray/blob/21202f6ddc3ceaf74fbc...

robertnishihara··on Amazon's exabyte-scale migration from Apache Spark to Ray on EC2
Your right that the serialization / deserialization overhead can quickly exceed the compute time. To avoid this you have to get a lot of small things right. And given our focus on ML workloads, this is particularly important when sharing large numerical arrays between processes (especially processes running on the same node).

One of the key things is to make sure the serialized object is stored in a data format where the serialized object does not need to be "transformed" in order to access it. For example, a numpy array can be created in O(1) time from a serialized blob by initializing a Python object with the right shape and dtype and a pointer to the right offset in the serialized blob. We also use projects like Apache Arrow that put a lot of care into this.

Example in more detail:

Imagine the object you are passing from process A to process B is a 1GB numpy array of floats. In the serialization step, process A produces a serialized blob of bytes that is basically just the 1GB numpy array plus a little bit of metadata. Process A writes that serialized blob into shared memory. This step of "writing into shared memory" still involves O(N) work, where N is the size of the array (though you can have multiple threads do the memcpy in parallel and be limited just by memory bandwidth).

In the deserialization step, process B accesses the same shared memory blob (process A and B are on the same machine). It reads a tiny bit of metadata to figure out the type of the serialized object and shape and so on. Then it constructs a numpy array with the correct shape and type and with a pointer to the actual data in shared memory at the right offset. Therefore it doesn't need to touch all of the bytes of data, it just does O(1) work instead of O(N).

That's the basic idea. You can imagine generalizing this beyond numpy arrays, but it's most effective for objects that include large numerical data (e.g., objects that include numpy arrays).

There are a bunch of little details to get right, e.g., serializing directly into shared memory instead of creating a serialized copy in process A and then copying it into shared memory. Doing the write into shared memory in parallel with a bunch of threads. Getting the deserialization right. You also have to make sure that the starting addresses of the numpy arrays are 64-byte aligned (if memory serves) so that you don't accidentally trigger a copy later on.

EDIT: I edited the above to add more detail.

robertnishihara··on Amazon's exabyte-scale migration from Apache Spark to Ray on EC2
I'm one of the creators of Ray. A few thoughts :)

1. This is truly impressive work from AWS. Patrick Ames began speaking about this a couple years ago, though at this point the blog post is probably the best reference. https://www.youtube.com/watch?v=h7svj_oAY14

2. This is not a "typical" Ray use case. I'm not aware of any other exabyte scale data processing workloads. Our bread and butter is ML workloads: training, inference, and unstructured data processing.

3. We have a data processing library called Ray Data for ingesting and processing data, often done in conjunction with training and inference. However, I believe in this particular use case, the heavy lifting is largely done with Ray's core APIs (tasks & actors), which are lower level and more flexible, which makes sense for highly custom use cases. Most Ray users use the Ray libraries (train, data, serve), but power users often use the Ray core APIs.

4. Since people often ask about data processing with Ray and Spark, Spark use cases tend to be more geared toward structured data and CPU processing. If you are joining a bunch of tables together or running SQL queries, Spark is going to be way better. If you're working with unstructured data (images, text, video, audio, etc), need mixed CPU & GPU compute, are doing deep learning and running inference, etc, then Ray is going to be much better.

robertnishihara··on Purple Llama: Towards open trust and safety in generative AI
We're hosting the model on Anyscale Endpoints. Try it out here [1]

[1] https://docs.endpoints.anyscale.com/supported-models/Meta-Ll...

robertnishihara··on JAX – NumPy on the CPU, GPU, and TPU
I'm a huge fan of Jax. The Jax team is incredibly strong!

Just want to share that Ray (an open source project we're developing at Anyscale), can be used to scale Jax (e.g., across TPUs).

Some docs from Google on how to do this

https://cloud.google.com/tpu/docs/ray-guide

Alpa is an open source project scaling Jax on 1000+ GPUs

https://www.anyscale.com/blog/training-175b-parameter-langua...

Cohere uses Ray + Jax + TPUs to build their LLMs

https://www.youtube.com/watch?v=For8yLkZP5w

A demo from Matt Johnson on the Jax team

https://www.youtube.com/watch?v=hyQ-tgD5sgc

robertnishihara··on A Comprehensive Guide for Building Rag-Based LLM Applications
Here is the blog post accompanying the notebook

https://www.anyscale.com/blog/a-comprehensive-guide-for-buil...

robertnishihara··on Fine-tune your own Llama 2 to replace GPT-3.5/4
It shouldn't be 100x. We've built an LLM API at Anyscale, and the price comparison works out as follows (per million tokens)

- Llama-2-70B: $1 (on Anyscale Endpoints [1]) - GPT-3.5-turbo: $1.50 - $2 (OpenAI [2])

[1] https://app.endpoints.anyscale.com/ [2] https://openai.com/pricing

robertnishihara··on Beating GPT-4 on HumanEval with a fine-tuned CodeLlama-34B
Thanks for the feedback, we'll improve the landing page!

The models (and current prices) right now are - Llama-2-7B ($0.25 / million tokens) - Llama-2-13B ($0.50 / million tokens) - Llama-2-70B ($1 / million tokens) - Code Llama ($1 / million tokens)

robertnishihara··on Beating GPT-4 on HumanEval with a fine-tuned CodeLlama-34B
It's amazing to see how rapidly things are moving.

You can try out CodeLlama-34B on Anyscale Endpoints (an LLM inference API we're building here at Anyscale for open source LLMs).

https://app.endpoints.anyscale.com/

robertnishihara··on Code Llama, a state-of-the-art large language model for coding
If you want to try out Code Llama, you can query it on Anyscale Endpoints (this is an LLM inference API we're working on here at Anyscale).

https://app.endpoints.anyscale.com/

robertnishihara··on GPT-3.5 Turbo fine-tuning and API updates
We've run experiments on datasets ranging from 5K - 100K examples, which gave fantastic results [1].

Some examples - https://huggingface.co/datasets/b-mc2/sql-create-context - https://huggingface.co/datasets/GEM/viggo

On the other hand, 8K examples was not enough to learn to solve grade school math problems [2], so it is very problem dependent.

[1] https://www.anyscale.com/blog/fine-tuning-llama-2-a-comprehe...

[2] https://huggingface.co/datasets/gsm8k

robertnishihara··on GPT-3.5 Turbo fine-tuning and API updates
I think for fine-tuned GPT-3.5 to be competitive with GPT-4 on your use cases (assistance with Angular), you'd have to fine-tune on enough data that it really resembles pre-training more than fine-tuning. And it wouldn't be worth the hassle unless you're building a product around it.

That said, many valuable LLM products / features are more narrow in scope and can see a huge lift from fine-tuning. We've run a bunch of experiments on this (e.g., SQL query generation is a good example), where fine-tuning even the 7B Llama-2 model outperforms GPT-4 (surprisingly) [1]. That's a very different type of problem from teaching software engineering of course.

[1] https://www.anyscale.com/blog/fine-tuning-llama-2-a-comprehe...

robertnishihara··on GPT-3.5 Turbo fine-tuning and API updates
If you want to query the Llama-2 models, you can use Anyscale Endpoints [1]. Note: I work on this :)

Llama-2-70B is $1 / million tokens, which is the most cost-efficient on the market that I'm aware of.

[1] https://app.endpoints.anyscale.com/

robertnishihara··on GPT-3.5 Turbo fine-tuning and API updates
I think of fine-tuning as an avenue to significantly reduce LLM inference costs, so I think this is an exciting development. You're right if you compare GPT-3.5-turbo to fine-tuned GPT-3.5-turbo, but if it's anything like fine-tuning the Llama-2 models, you'll be able to achieve GPT-4 level performance for a wide range of practical use cases (SQL query generation is an example), but probably not for math or coding (at least not without fine-tuning on a significant amount of data).

In fact, we've seen GPT-4 level performance from even the 7B Llama-2 model after fine-tuning. [1]

[1] https://www.anyscale.com/blog/fine-tuning-llama-2-a-comprehe...

robertnishihara··on Azure ChatGPT: Private and secure ChatGPT for internal enterprise use
> Llama 2 might by some measures be close to GPT 3.5, but it’s nowhere near GPT 4

I think you're right about this, and benchmarks we've run at Anyscale support this conclusion [1].

The caveat there (which I think will be a big boon for open models) is that techniques like fine-tuning makes a HUGE difference and can bridge the quality gap between Llama-2 and GPT-4 for many (but not all) problems.

[1] https://www.anyscale.com/blog/fine-tuning-llama-2-a-comprehe...

robertnishihara··on Azure ChatGPT: Private and secure ChatGPT for internal enterprise use
If you want to try out the Llama-2 models (7B, 13B, 70B), you can get started very easily with Anyscale Endpoints (~2 min). https://app.endpoints.anyscale.com/
robertnishihara··on Azure ChatGPT: Private and secure ChatGPT for internal enterprise use
> we believe most companies prefer locally installed solutions to cloud based ones

We've also seen a strong desire from businesses to manage models and compute on their own machines or in their own cloud accounts. This is often part of a hybrid strategy of using API products like OpenAI for rapid prototyping.

The majority of (though not all) businesses we've seen tend to be quite comfortable using hosted API products for rapid prototyping and for proving out an initial version of their AI functionality. But in many cases, they want to complement that with the ability to manage models and compute themselves. The motivation here is often to reduce costs by using smaller / faster / cheaper fine-tuned open models.

When we started Anyscale, customer demand led us to run training & inference workloads in our customers' cloud accounts. That way your data and code stays inside of your own cloud account.

Now with all the progress in open models and the desire to rapidly prototype, we're complementing that with a fully-managed inference API where you can do inference with the Llama-2 models [1] (like the OpenAI API but for open models).

[1] https://app.endpoints.anyscale.com/

robertnishihara··on Azure ChatGPT: Private and secure ChatGPT for internal enterprise use
We (at Anyscale) have benchmarked GPT-4 versus the Llama-2 suite of models on a few problems: functional representation, SQL generation, grade-school math question answering.

GPT-4 wins by a lot out of the box. However, surprisingly, fine-tuning makes a huge difference and allows the 7B Llama-2 model to outperform GPT-4 on some (but not all) problems.

This is really great news for open models as many applications will benefit from smaller, faster, and cheaper fine-tuned models rather than a single large, slow, general-purpose model (Llama-2-7B is something like 2% of the size of GPT-4).

GPT-4 continues to outperform even the fine-tuned 70B model on grade-school math question answering, likely due to the data Llama-2 was trained on (more data for fine-tuning helps here).

https://www.anyscale.com/blog/fine-tuning-llama-2-a-comprehe...

robertnishihara··on Launch HN: BuildFlow (YC W23) – The FastAPI of data pipelines
Congratulations on the launch, I love the focus on ease of use and making it easy to get started, and it's exciting to see impressive products being built with Ray!

I'm one of the Ray developers. It is true that Ray focuses a lot on ML applications (in particular, the main libraries built on top of Ray are for workloads like training, serving, and batch processing / inference). That said, one of our long-term goals with Ray is to be a great general-purpose way to build distributed applications, so I hope it is working out for you :)

robertnishihara··on How to train large models on many GPUs? (2021)
Yes, but this will largely come down to whether the deep learning framework that you're using (PyTorch, TensorFlow, Jax, etc) works well in that setting. Ray is pretty framework and hardware agnostic and can be used to schedule / scale different ML frameworks on different types of devices (CPUs, GPUs, TPUs, etc), but the actual logic for running code on the accelerators lives in the deep learning framework.
robertnishihara··on How to train large models on many GPUs? (2021)
I'm one of the Ray developers, thanks for the shoutout :)

If you're curious about how Ray is used for LLMs, here are some interesting examples of LLM projects using Ray!

- Alpa does training and serving with 175B parameter models https://github.com/alpa-projects/alpa

- GPT-J https://github.com/kingoflolz/mesh-transformer-jax

- Another HN thread on training LLMs with Ray (on TPUs in this case) https://news.ycombinator.com/item?id=27731168

- OpenAI fireside chat on the evolution of their infrastructure and usage of Ray for training https://www.youtube.com/watch?v=CqiL5QQnN64

- Cohere on their architecture for training LLMs https://www.youtube.com/watch?v=For8yLkZP5w&t=3s

Some other thoughts

1. There is a lot more we want to do to make Ray better for working with large language models and for making training, serving, and batch inference work well out of the box.

2. The original post is about training, but we actually see even more interest in fine-tuning and serving with LLMs, in part because there are good pre-trained models.

3. For LLMs, we see a lot of interest in Ray + Jax or Ray + TPUs relative to what we see in other use cases.

robertnishihara··on Ray: A Distributed Framework for Emerging AI Applications
Thanks for the comments! A few quick notes

The term serverless is a bit overloaded. Here's what we want. (1) Ray users should be able to focus only on their application logic and not have to worry about configuring clusters at least in the common case (this is not saying that they can't configure clusters if they want to). (2) In addition, we want Ray applications to be portable and to run on different size clusters with different instance types on different cloud providers or k8s or your laptop or anywhere. The application should be decoupled from the cluster configuration. (3) Great support for autoscaling Ray clusters. When a Ray application needs more resources of a certain type (CPUs, GPUs, memory, etc), those resources should be added to the cluster. When they are no longer needed, the cluster should scale down.

This is quite different from FaaS, though seems in line with the spirit of serverless. And these "serverless" properties are not something we think of as separate from Ray, but rather as part of just making Ray work better.

robertnishihara··on Ray: A Distributed Framework for Emerging AI Applications
Hi all, I'm one of the authors of Ray, thanks for all the comments and discussion! To add to the discussion, I'll mention a few conceptual things that have changed since we wrote the paper.

*Emphasis on the library ecosystem*

A lot of our focus is on building an ecosystem of libraries on top of Ray (much, but not all, of the focus is on machine learning libraries). Some of these libraries are built natively on top of Ray such as Ray Tune for scaling hyperparameter search (http://tune.io), RLlib for scaling reinforcement learning (http://rllib.io), Ray Serve for scaling model serving (http://rayserve.org/), and RaySGD for scaling training (https://docs.ray.io/en/master/raysgd/raysgd.html).

Some of the libraries are popular libraries on their own, which now integrate with Ray such as Horovod (https://eng.uber.com/horovod-ray/), XGBoost (https://xgboost.readthedocs.io/en/latest/tutorials/ray.html), and Dask for dataframes (https://docs.ray.io/en/master/dask-on-ray.html). While Dask itself has similarities to Ray (especially the task part of the Ray API), Dask also has libraries for scaling dataframes and arrays, which can be used as part of the Ray ecosystem (more details at https://www.anyscale.com/blog/analyzing-memory-management-an...).

Many Ray users start using Ray for one of the libraries (e.g., to scale training or hyperparameter search) as opposed to just for the core system.

*Emphasis on serverless*

Our goal with Ray is to make distributed computing as easy as possible. To do that, we think the serverless direction, which allows people to just focus on their code and not on infrastructure, is very important. Here, I don't mean serverless purely in the sense of functions as a service, but something that would allow people to run a wide variety of applications (training, data processing, inference, etc) elastically in the cloud without configuring or thinking about infrastructure. There's a lot of ongoing work here (e.g., to improve autoscaling up and down with heterogeneous resource types). More details on the topic https://www.anyscale.com/blog/the-ideal-foundation-for-a-gen....

If you're interested in this kind of stuff, consider joining us at Anyscale https://jobs.lever.co/anyscale.