This is not a cut and dry question. A single node cluster gives you the lowest latency because all your data will reside in your single vm along with your application because you don't incur any additional costs in serialization and chattiness. Obviously, this defeats the purpose of Hazelcast which provides for you the features required for the CAP Theorem like High Availability. That aside, you simply may not fit all your data in a single node. From here, as you add more nodes, that latency may increase if your embedded app is accessing data from another node, especially if it has multiple round-trips. But, you will more likely increase throughput, thus the ability to process more concurrently, because you'll have more cpu and memory available across all your nodes.
Sizing your cluster will depend on the usual variables such as: available heap size, data (object) size, count, indexes, retention policy, and estimated growth. There are many other variables to be considered that will affect the size of your cluster.