I feel like Erlang/Elixir is designed to handle cases like this robustly, but it's not clear to me how this code avoids losing data when a process crashes, potentially up to 5s worth of updates!
I feel like Erlang/Elixir is designed to handle cases like this robustly, but it's not clear to me how this code avoids losing data when a process crashes, potentially up to 5s worth of updates!
Long Answer - What happens in any language where you batch stuff in memory? If the system breaks down, you lose that data. This isn't unique to Erlang/Elixir. There's a tradeoff that persistent storage is slower than memory, but it's persistent. So do you want to be fast, or do you want to be durable?
However, there is some flexibility to address this at a system level. Namely, you wait for confirmation. A client, be that an actual user, another process, etc, wants to know if the write has been persisted. Whereas many other languages make this sort of caching layer completely transparent to the client (i.e., they return success immediately, and so failures mean invisible data loss as the cache is dropped), Erlang/Elixir's model makes it so you HAVE to think about this. I sent a message to the downstream process; do I care about a response? If I don't get a response for any given interval, I don't know what happened to that message. I can still do useful work in the meantime, but I don't know that my message was fully handled (persisted). I can retry or I can report failure or whatever.
This is true in any distributed system, and in Erlang/Elixir it's expressed very evidently in the language constructs (rather than being hidden from you).
You can build local disk caches if you want (in fact, there are included tools to make this super easy for you; ETS, DETS, and Mnesia allow you to shove Erlang terms into memory storage, disk based storage, and a hybrid of the two with some nice DB-like behaviors, respectively), but you need to choose to slow your message ingestion to the speed of local disk writes in that case (as well as handle synchronization of deletes in the event of multiple writers to your downstream). An depending what your upstream is, that doesn't provide a guarantee (i.e., a write to local disk != a persisted write from a user perspective, because the disk could crash before it ever makes it off the local one).
A little over a decade ago I did a project using Erlang. My next project was using Ruby, and I was suddenly _horrified_ by the fact that I had no idea what would happen if the application crashed (and Ruby apps, at least in those days, crashed fairly often). Erlang/Elixir both forces you to think about these things, and gives you tools to address them, where many other languages (or their libraries) simply assume that we will stay on the happy path.
Of course, you can be very productive using languages that just ignore the possibility of very rare failures. Many successful systems work on the basis that sometimes shit just happens and maybe some data does get lost. But, after programming with Erlang or Elixir for a while, that situation starts to feel less acceptable!
I consider the couple years I built systems in Erlang to be fundamental for me. It's affected, for better and worse, my entire approach to system design, at every level. It's meant the stuff I or (now that I'm in management) my teams tend to write is incredibly resilient (compared with the other teams in the department), but also meant that I have a really hard time with any Silicon Valley interview.
If the process dies, the Supervisor invokes the `terminate/2` callback so you're still able to process events.
Here are the relevant lines: https://github.com/plausible/analytics/blob/b724def948d51a0f...
Keep in mind you need to trap exit signal to tell the supervisor to invoke the callback, as done so here: https://github.com/plausible/analytics/blob/b724def948d51a0f...
The erlang docs also mentions this: https://erlang.org/doc/design_principles/gen_server_concepts...
If I have as significantly long transaction, I can restart the transaction from transaction_id + 1 forward and then handle the actual transaction out of band by hand if its something very odd.
I usually insert a logger in these code paths too, but the crash handler is usually after the logger, which makes it easier to reason what has happened if you are dealing with mutable state.
https://elixir-lang.org/getting-started/mix-otp/supervisor-a...
Can someone chime in on this?
Calls by contrast check out a monitor on the counterparty and crash with a timeout so you have guaranteed delivery and acknowledgement, or crash the calling process.
A lot of times people from other plarforms rush to use casts when they "don't need a response" but the actual meanings have more to do with failure domains and rate limiting back pressure; my personal feeling is you should default to call and only use cast when you need failure isolation.
If you use erlang:send/3 with options or erlang:send_nosuspend/2, you can have some cases where messages would be dropped without trying. I think that may be the source of the drop under load you're thinking of?
Without nosuspend, if you send to a node that's connected, dist will queue it to be sent, but there's no guarantee it's received because networking, and it may not even be sent if the other side has gone away and the send queue is already too large; There is a hard to express constraint that if your message doesn't get received in this case, dist will disconnect eventually, but it's important to note that the tick time-out doesn't guarantee timely delivery either; as long as some data is flowing, you can get quite a backlog; I think I've gotten net_adm:ping times above 30 minutes in some cases. Also, if the dist connection is dropped, but both nodes are online, it will likely reconnect shortly, and some messages will have been lost.
If you send to a node that's not connected, and didn't specify no_connect, dist will queue the message while attempting to connect, but if that attempt fails, the messages will be dropped.
It's also possible for code running as a gen_server (or GenServer, I suppose) to check how many messages are queued for it, and run different logic. I've written gen_servers that would drop optional requests if the queue was large. Also, if client timeouts are known, there are ways to approximate the time spent waiting in queue, and drop requests if they are received after the client already timed out; it's a little tricky to do this though.
I'll chime in explicitly to your question though - messages don't just disappear without a reason, but the reason may not be visible to you. It's possible a receiving process got a message, crashed partway through processing it (possibly after you even sent other messages, so dropping those, too), restarted, and now is handling the next incoming messages, making it look like a message dropped. It's possible a receiving process is on another node (in distributed Erlang) and the packet dropped. Etc etc.
These all look the same. The way to handle them is always the same; accept you have at most once delivery, or include an 'ack or retry' mechanism and write your system to handle at least once delivery. As that article mentions, Erlang makes you start thinking of your system as a distributed system from the get go. This initially feels incredibly inconvenient, but as you go down the distributed system, fault tolerant, CAP-bound rabbit hole, it becomes incredibly helpful. The system doesn't hide the things it can't guarantee from you; it forces you to feel uneasy about them.
(Whereas, for example, in Java, I don't know that that worker thread has terminated. So some other thread's fiddling with some bit of shared memory to communicate something to that thread does nothing. Which I've had happen in production and led to a single instance out of the fleet have stale data, which caused us no end of grief. Erlang forces you to ask, as you've just done, "what happens if this message doesn't get received because something happened to the receiver"?)
So for example one possible implementation of state is that you catch your own crash, save some state and then restore it when you restart
You aren't supposed to ever crash, however the whole language it's based around designing a framework to handle what happens if you do crash
So you design a whole framework of supervisors, watching processes, that all have a defined startup and shutdown order. You build structures with rules such as: if my cache module fails, them restart this whole chunk of application over here.
it's very cool
Obviously, different tolerances to that are going to necessitate different designs. You can configure the size of the message buffer when a GenServer starts, and if you have very low tolerance for lost data, you'd want to use a synchronous message (using call instead of cast, which blocks the sender until it returns) and appropriate error handling.
BEAM has a plethora of features for reliable applications, you just have to apply them appropriately.