[ 
https://issues.apache.org/jira/browse/SPARK-59621?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
 ]

ASF GitHub Bot updated SPARK-59621:
-----------------------------------
    Labels: pull-request-available  (was: )

> Serialize StatusUpdate manually (Externalizable) to avoid 
> Enumeration/BigDecimal overhead
> -----------------------------------------------------------------------------------------
>
>                 Key: SPARK-59621
>                 URL: https://issues.apache.org/jira/browse/SPARK-59621
>             Project: Spark
>          Issue Type: Improvement
>          Components: Spark Core
>    Affects Versions: 4.1.0
>            Reporter: David Mollitor
>            Priority: Minor
>              Labels: pull-request-available
>
> h3. Description
> {{CoarseGrainedClusterMessages.StatusUpdate}} is sent from executor to driver 
> on *every task state change* ({{{}RUNNING{}}} at start, 
> {{{}FINISHED{}}}/{{{}FAILED{}}}/{{{}KILLED{}}} at end) – roughly two or more 
> per task, on the driver's RPC intake path.
> Measuring its Java-serialized size (a fresh {{ObjectOutputStream}} per RPC, 
> so no class-descriptor caching) shows an empty-payload {{StatusUpdate}} is 
> {*}1713 bytes{*}, and two fields account for ~75% of it:
> ||Field||Serialized bytes||
> |{{state}} (a Scala {{Enumeration}} value)|642|
> |{{taskCpus}} (a {{{}BigDecimal{}}})|638|
> |{{resources}} (empty {{{}Map{}}})|150|
> |{{data}} (empty {{{}SerializableBuffer{}}})|68|
> |full {{StatusUpdate}}|1713|
> A Scala {{Enumeration.Value}} serializes a reference to its enclosing 
> {{Enumeration}} object, so it drags in the whole {{TaskState}} object; 
> {{scala.math.BigDecimal}} drags in
> {{java.math.BigDecimal}} + {{BigInteger}} + {{MathContext}} + 
> {{RoundingMode}} descriptors. The actual payload is ~31 bytes.
> Make {{StatusUpdate}} {{Externalizable}} with a compact manual encoding, 
> mirroring
> {{UpdateBlockInfo}} in the same message family and reusing 
> {{{}TaskDescription{}}}'s exact wire form for the fractional-CPU 
> {{BigDecimal}} ({{{}CpuAmount.toDisplayString{}}} / 
> {{{}CpuAmount.normalize{}}}): {{state}} as one byte ({{{}.id{}}}), 
> {{taskCpus}} as its normalized decimal string, {{data}} via 
> {{{}SerializableBuffer{}}}'s existing channel-based serialization (no extra 
> copy for large results), and {{resources}} as a size-prefixed nested map.
> This shrinks the empty-payload message from *1713 bytes to 191 bytes* (~9x 
> smaller), cutting steady serialization/GC/bandwidth on a per-task-frequency 
> control message. It is behavior-preserving (all fields round-trip; fractional 
> CPUs exactly).
> Wire-compat: this is an internal driver<->executor RPC within a single 
> application, and Java serialization is already documented as not stable 
> across Spark versions.



--
This message was sent by Atlassian Jira
(v8.20.10#820010)

---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to