Split and merge records
Split and Merge algorithm is very helpful to accelerate execution of a Business Process by running tasks in parallel rather than sequentially. For example, if a document consists of several pages, we can split them, process and then merge back to work with a single document from this point.
Split and merge records via rule
Few details to consider in this approach:
- Business Process streaming option should be switched off (Run Business Process > Advanced Options > Threshold = 100% ). This means each step of BP will 'wait' until all tasks will be processed, only after that next step starts to process.
- Grouping by Rule contains hard-coded value for column name, see example below.
- Flow can work with any amount of initial data, no restrictions (Option 2 works with 1 record on Input).
Usage:
- Create first Bot Task according to the sample, where the 'group_by_column' is a mandatory column to group by all initial data and is used in the grouping rule
Expand to see the sample
<config>
<var-def name="group_by_column">
<template>${UUID.randomUUID()}</template>
</var-def>
<script></script>
<while condition="true" maxloops="10" empty="empty">
<script></script>
</while>
<export include-original-data="false">
<multi-column list="${exportItems}" split-results="true">
<put-to-column-getter name="group_by_column" property="group_by_column" />
<put-to-column-getter name="column_value_1" property="column_value_1" />
<put-to-column-getter name="column_value_2" property="column_value_2" />
</multi-column>
</export>
</config>
- Grouping by rule
Expand to see the sample
package com.freedomoss.requester;
import com.freedomoss.objective.model.CompositeRuleContext;
import com.freedomoss.objective.model.CompositeRuleDataItemContext;
import org.slf4j.Logger;
import java.util.Map;
import java.util.HashMap;
import java.util.Collections;
import java.util.LinkedHashMap;
import com.freedomoss.requester.model.AwsHitQuestionItem;
import java.util.ArrayList;
import java.util.List;
import java.util.Iterator;
import com.freedomoss.requester.dto.SubmissionResultDTO;
import java.util.UUID;
global CompositeRuleContext source
global Logger log
global Map params
# RuleContext.JOIN_RULE
# RuleContext.USE_ONLY_LAST_RUN_DATA
rule "Rule context initialization"
auto-focus true
dialect "mvel"
no-loop
agenda-group "initialization-group"
when
$ctx:CompositeRuleContext(initialized == false);
then
$ctx.properties["DISPLAY_INDEX_MAP"] = new LinkedHashMap();
# insert processed facts into memory
$ctx.updateWorkingMemory();
# move to business rules
kcontext.getKnowledgeRuntime().getAgenda().getAgendaGroup("calculation").setFocus();
end
rule "1. Grouping records"
agenda-group "calculation"
no-loop
dialect "mvel"
salience 200
when
$ctx:CompositeRuleContext(initialized == true, $props:properties);
$rqc:CompositeRuleDataItemContext(this memberOf $ctx.lastStepDataItems);
then
log.info("Grouping records ");
accum($props, $rqc, log);
end
rule "2. Submitting records to the next step"
agenda-group "calculation"
no-loop
dialect "mvel"
salience 150
when
$ctx:CompositeRuleContext(initialized == true);
then
log.info("Submitting records to the next step ");
submit($ctx, log);
$ctx.logExecutedRule(kcontext.getRule().getName());
end
function void submit(CompositeRuleContext ctx, Logger log) {
Map map = (Map) ctx.getProperties().get("DISPLAY_INDEX_MAP");
log.info("Submitting records, " + map.keySet().size());
for (Iterator iterator = map.keySet().iterator(); iterator.hasNext();) {
String uKey = (String) iterator.next();
log.info("uKey=" + uKey);
SubmissionResultDTO nq = new SubmissionResultDTO(UUID.randomUUID().toString(), null);
nq.putData((Map) map.get(uKey));
ctx.sendResultTo("original", nq);
}
}
function void accum(Map props, CompositeRuleDataItemContext item, Logger log) {
Map map = (Map) props.get("DISPLAY_INDEX_MAP");
log.info("item::: " + item);
String GROUPING_FIELD = "group_by_column";
Map itemValues = item.getAllValues();
String grouping_field_value = (String) itemValues.get(GROUPING_FIELD);
log.info(GROUPING_FIELD + ": " + grouping_field_value);
if (!map.containsKey(grouping_field_value)) {
log.info("Creating new stored map");
Map newStoredMap = new LinkedHashMap();
List columns_ = new ArrayList(itemValues.keySet());
Collections.sort(columns_);
for (Iterator listIterator = columns_.iterator(); listIterator.hasNext();) {
String key = (String) listIterator.next();
String storedValue = (String) itemValues.get(key);
newStoredMap.put(key, storedValue);
}
map.put(grouping_field_value ,newStoredMap);
}
}

Input Data: group_by.txt
Package to import: Split+and+join+sample+TEST+23-4-2019.zip
Some more examples for split join rules:
We need to merge (group) results after join rule. Results should be joined by columns and in each column values (all values from all records) should be separated by pipes (|). The BP diagram looks as follows.

Package to import: Split+and+join+sample+TEST_v10+24-4-2019.zip
Explicit grouping by rule: new-subm-accum-split-rule.txt
Create multiple records and split them using original data
Business case:
- We have 1 record at the start.
- We want to keep original record from start and split it to multiple records.
- All new records will have original data plus new data.
Expand to see the example
<config>
<var-def name="multiple_dividends">
<template>${multiple_dividends}</template>
</var-def>
<script><![CDATA[
List itemsForExport = new ArrayList();
List columns = new ArrayList();
List values = hit_submission_data_item.getWrappedObject().getItemValueList();
int originalListSize = values.size();
for(int i = 0; i < originalListSize; i++) {
com.freedomoss.crowdcontrol.webharvest.HitSubmissionDataItemValueDto dataItem = (com.freedomoss.crowdcontrol.webharvest.HitSubmissionDataItemValueDto) values.get(i);
String name = dataItem.getName();
columns.add(name);
}
columns.add("payment_date");
columns.add("issuer");
columns.add("record_date");
columns.add("gross_amount");
columns.add("stock_type");
columns.add("currency");
columns.add("ex_date");
columns.add("dividend_type");
columns.add("tax_rate");
columns.add("frank_rate");
columns.add("unfranked_amount");
columns.add("franked_amount");
columns.add("drp_price");
columns.add("conduit_foreign_income");
columns.add("election_date");
columns.add("meeting_date");
columns.add("conditional");
columns.add("conditional_criteria");
columns.add("announcement_date");
columns.add("unfranked_tax_rate");
//columns.add("interest");
//columns.add("selic_rate");
//columns.add("conversion_rate");
]]>
</script>
<case>
<if condition='${!multiple_dividends.toString().isEmpty() &&multiple_dividends.toString().contains("{") }'>
<var-def name="multiple_divdends_json_val">
<json-to-xml>
<template>${multiple_dividends}</template>
</json-to-xml>
</var-def>
<var-def name="results">
<xpath expression="//multiple_dividends">
<template><content>${multiple_divdends_json_val}</content></template>
</xpath>
</var-def>
<loop item="item">
<list>
<var name="results" />
</list>
<body>
<empty>
<var-def name="payment_date_value_ret">
<xpath expression="//payment_date_value/text()">
<var name="item" />
</xpath>
</var-def>
<var-def name="record_date_value_ret">
<xpath expression="//record_date_value/text()">
<var name="item" />
</xpath>
</var-def>
<var-def name="gross_amount_value_ret">
<xpath expression="//gross_amount_value/text()">
<var name="item" />
</xpath>
</var-def>
<var-def name="stock_type_value_ret">
<xpath expression="//stock_type_value/text()">
<var name="item" />
</xpath>
</var-def>
<var-def name="currency_value_ret">
<xpath expression="//currency_value/text()">
<var name="item" />
</xpath>
</var-def>
<var-def name="ex_date_value_ret">
<xpath expression="//ex_date_value/text()">
<var name="item" />
</xpath>
</var-def>
<var-def name="dividend_type_value_ret">
<xpath expression="//dividend_type_value/text()">
<var name="item" />
</xpath>
</var-def>
<var-def name="meeting_date_value_ret">
<xpath expression="//meeting_date_value/text()">
<var name="item" />
</xpath>
</var-def>
<var-def name="tax_rate_value_ret">
<xpath expression="//tax_rate_value/text()">
<var name="item" />
</xpath>
</var-def>
<var-def name="conditional_value_ret">
<xpath expression="//conditional_value/text()">
<var name="item" />
</xpath>
</var-def>
<var-def name="announcement_date_value_ret">
<xpath expression="//announcement_date_value/text()">
<var name="item" />
</xpath>
</var-def>
<var-def name="issuer_value_ret">
<xpath expression="//issuer_value/text()">
<var name="item" />
</xpath>
</var-def>
<var-def name="frank_rate_value_ret">
<xpath expression="//frank_rate_value/text()">
<var name="item" />
</xpath>
</var-def>
<var-def name="unfranked_amount_value_ret">
<xpath expression="//unfranked_amount_value/text()">
<var name="item" />
</xpath>
</var-def>
<var-def name="franked_amount_value_ret">
<xpath expression="//franked_amount_value/text()">
<var name="item" />
</xpath>
</var-def>
<var-def name="unfranked_tax_rate_value_ret">
<xpath expression="//unfranked_tax_rate_value/text()">
<var name="item" />
</xpath>
</var-def>
<var-def name="election_date_value_ret">
<xpath expression="//election_date_value/text()">
<var name="item" />
</xpath>
</var-def>
<var-def name="drp_price_value_ret">
<xpath expression="//drp_price_value/text()">
<var name="item" />
</xpath>
</var-def>
<var-def name="conduit_foreign_income_value_ret">
<xpath expression="//conduit_foreign_income_value/text()">
<var name="item" />
</xpath>
</var-def>
<var-def name="total_dividend_amount_ret">
<xpath expression="//total_dividend_amount_value/text()">
<var name="item" />
</xpath>
</var-def>
<var-def name="total_dividend_amount_estimated_or_actual_ret">
<xpath expression="//total_dividend_amount_estimated_or_actual_value/text()">
<var name="item" />
</xpath>
</var-def>
<!--<var-def name="interest_value_ret">
<xpath expression="//interest_value/text()">
<var name="item" />
</xpath>
</var-def>
<var-def name="selic_rate_value_ret">
<xpath expression="//selic_rate_value/text()">
<var name="item" />
</xpath>
</var-def>
<var-def name="conversion_rate_value_ret">
<xpath expression="//conversion_rate_value/text()">
<var name="item" />
</xpath>
</var-def> -->
<script><![CDATA[
Map rec = new HashMap();
for(int i = 0; i < originalListSize; i++) {
com.freedomoss.crowdcontrol.webharvest.HitSubmissionDataItemValueDto dataItem = (com.freedomoss.crowdcontrol.webharvest.HitSubmissionDataItemValueDto) values.get(i);
String name = dataItem.getName();
String value = dataItem.getValue();
rec.put(name, value);
}
rec.put("payment_date", payment_date_value_ret.toString().trim());
rec.put("record_date", record_date_value_ret.toString().trim());
rec.put("gross_amount", gross_amount_value_ret.toString().trim());
rec.put("stock_type", stock_type_value_ret.toString().trim().replace("|",""));
rec.put("currency", currency_value_ret.toString().trim());
rec.put("ex_date", ex_date_value_ret.toString().trim());
rec.put("dividend_type", dividend_type_value_ret.toString().trim());
rec.put("tax_rate", tax_rate_value_ret.toString().trim());
rec.put("meeting_date", meeting_date_value_ret.toString().trim());
if (conditional_value_ret.toString().trim().isEmpty()) {
rec.put("conditional", "No");
}
else {
rec.put("conditional", "Yes");
}
rec.put("conditional_criteria", conditional_value_ret.toString().trim());
rec.put("issuer", issuer_value_ret.toString().trim());
rec.put("frank_rate", frank_rate_value_ret.toString().trim());
rec.put("announcement_date", announcement_date_value_ret.toString().trim());
rec.put("unfranked_amount", unfranked_amount_value_ret.toString().trim());
rec.put("franked_amount", franked_amount_value_ret.toString().trim());
rec.put("unfranked_tax_rate", unfranked_tax_rate_value_ret.toString().trim());
rec.put("drp_price", drp_price_value_ret.toString().trim());
rec.put("conduit_foreign_income", conduit_foreign_income_value_ret.toString().trim());
rec.put("election_date", election_date_value_ret.toString().trim());
rec.put("total_dividend_amount", total_dividend_amount_ret.toString().trim());
rec.put("total_dividend_amount_estimated_or_actual", total_dividend_amount_estimated_or_actual_ret.toString().trim());
//rec.put("interest", interest_value_ret.toString().trim());
//rec.put("selic_rate", selic_rate_value_ret.toString().trim());
//rec.put("conversion_rate", conversion_rate_value_ret.toString().trim());
itemsForExport.add(rec);
]]></script>
</empty>
</body>
</loop>
</if>
</case>
<export include-original-data="false">
<multi-column list="${itemsForExport}" split-results="true">
<loop item="columnName">
<list>
<script return="columns" />
</list>
<body>
<put-to-column-getter name="${columnName}" property="${columnName}" />
</body>
</loop>
</multi-column>
</export>
</config>