Yes, I am caching; there's definitely garbage collection going on, and I am
pushing the limits of our memory, though I think we still fit; thanks, I'll
take a look at the wait parameter.


On Mon, Apr 14, 2014 at 8:39 PM, Matei Zaharia <[email protected]>wrote:

> Spark can actually launch multiple executors on the same node if you
> configure it that way, but if you haven’t done that, this might mean that
> some tasks are reading data from the cache, and some from HDFS. (In the
> HDFS case Spark will only report it as NODE_LOCAL since HDFS isn’t tied to
> a particular executor process). For example, maybe you cached some data but
> not all the partitions of the RDD are in memory. Are you using caching here?
>
> There’s a locality wait setting in Spark (spark.locality.wait) that
> determines how long it will wait to go to the next locality level when it
> can’t launch stuff at its preferred one (e.g. to go from process to node).
> You can try increasing that too, by default it’s only 3000 ms. It might be
> that the whole RDD is cached but garbage collection causes it to give up
> waiting on some nodes and launch stuff on other nodes instead, which might
> be HDFS-local (due to data replication) but not cache-local.
>
> Matei
>
> On Apr 14, 2014, at 8:37 AM, dachuan <[email protected]> wrote:
>
> > I am confused about the process local and node local, too.
> >
> > In my current understanding of Spark, one application typically only has
> one executor in one node. However, node local means your data is in the
> same host, but in a different executor.
> >
> > This further means node local is the same with process local unless one
> node has two executors, which could only happen when one node has two
> Workers.
> >
> > Waiting for further discussion ..
> >
> >
> > On Mon, Apr 14, 2014 at 10:13 AM, Nathan Kronenfeld <
> [email protected]> wrote:
> > I've a fairly large job (5E9 records, ~1600 partitions).wherein on a
> given stage, it looks like for the first half of the tasks, everything runs
> in process_local mode in ~10s/partition.  Then, from halfway through,
> everything starts running in node_local mode, and takes 10x as long or more.
> >
> > I read somewhere that the difference between the two had to do with the
> data being local to the running jvm, or another jvm on the same machine.
>  If that's the case, shouldn't the distribution of the two modes be more
> random?  If not, what exactly is the difference between the two modes?
>  Given how much longer it takes in node_local mode, it seems like the whole
> thing would probably run much faster just by waiting for the right jvm to
> be free.  Is there any way of forcing this?
> >
> >
> > Thanks,
> >               -Nathan
> >
> >
> > --
> > Nathan Kronenfeld
> > Senior Visualization Developer
> > Oculus Info Inc
> > 2 Berkeley Street, Suite 600,
> > Toronto, Ontario M5A 4J5
> > Phone:  +1-416-203-3003 x 238
> > Email:  [email protected]
> >
> >
> >
> > --
> > Dachuan Huang
> > Cellphone: 614-390-7234
> > 2015 Neil Avenue
> > Ohio State University
> > Columbus, Ohio
> > U.S.A.
> > 43210
>
>


-- 
Nathan Kronenfeld
Senior Visualization Developer
Oculus Info Inc
2 Berkeley Street, Suite 600,
Toronto, Ontario M5A 4J5
Phone:  +1-416-203-3003 x 238
Email:  [email protected]

Reply via email to