Skip to main content
The Parallel Bulk Loader (PBL) job enables bulk ingestion of structured and semi-structured data from big data systems, NoSQL databases, and common file formats like Parquet and Avro. Datasources the PBL uses include not only common file formats, but Solr databases, JDBC-compliant databases, MongoDB databases and more. In addition, the PBL distributes the load across the Lucidworks Search Spark cluster to optimize performance. And because no parsing is needed, indexing performance is also maximized by writing directly to Solr.

Usage

Use this job to load data into Lucidworks Search from a SparkSQL-compliant datasource, and then send the data to any Spark-supported datasource such as Solr, an index pipeline, etc. To create a Parallel Bulk Loader job, sign in to Lucidworks Search and click Collections > Jobs. Then click Add+ and in the Custom and Others Jobs section, select Parallel Bulk Loader. You can enter basic and advanced parameters to configure the job. If the field has a default value, it is populated when you click to add the job. Parallel Bulk Loader can be configured for many different use cases. Two examples are:
  • Organizations that need to meet financial/bank-level transactional integrity requirements use SQL databases. Those datasources are structured in relational tables and ensure data integrity even if errors or power failures occur. In addition, these datasources are useful in large-scale operations that employ complex queries for analytics and reporting. For example, if categories and products have various prices, discounts, and offering dates, a SQL datasource is the most efficient option. Lucidworks supports SQL databases such as JDBC-compliant databases.
  • In contrast, NoSQL databases are based on documents and allow for more flexibility with structured, semi-structured, and unstructured data. For example, your organization might need information about user product reviews or session data. And if your organization needs to process massive amount of data from multiple systems, NoSQL is an efficient option. Lucidworks supports NoSQL databases such as MongoDB and Apache Cassandra.

Example configuration

The following is an example configuration. The table after the configuration defines the fields.

Additional field information

This section provides more detailed information about some of the configuration fields.
  • format. Spark scans the job’s classpath for a class named DefaultSource in the <format> package. For the solr format, where the solr.DefaultSource class is defined in the spark-solr repository.
  • transformScala. If the Scala script is not sufficient, you might need the full power of the Spark API to transform data into an indexable form.
    • The transformScala option lets you filter and/or transform the input DataFrame any way you would like. You can even define UDFs to use during your transformation. For an example of using Scala to transform the input DataFrame before indexing in Solr, see the Import with the Bulk Loader example.
    • Another powerful use of the transformScala option is that you can pull in advanced libraries, such as Spark NLP (from John Snow Labs) to do NLP work on your content before indexing. See the Import with the Bulk Loader example.
    • Your Scala script can do other things but, at a minimum, it must define the following function that the Parallel Bulk Loader invokes:
  • cacheAfterRead. This hidden field specifies if input data is retained in memory (and on disk as needed) after reading. If set to true, it may help stability of the job by reading all data from the input source first before transforming or writing to Solr. This could make the job run slower because it adds an intermediate write operation. For example, false.
  • updates. This field lists the userId accounts that have been updated. The timestamp contains the date and time the account was updated. The value is displayed in Unix epoch time in a yyyy-mm-ddThh:mm:ssZ format. For example, the power_user was updated on 2024-07-30T20:30:31.350243642Z.

Create and run Parallel Bulk Loader jobs

Use the Jobs manager to create and run Parallel Bulk Loader jobs. You can also use the Scheduler to schedule jobs.In the procedures, select Parallel Bulk Loader as the job type and configure the job using the following information.

Configuration settings for the Parallel Bulk Loader job

This section provides configuration settings for the Parallel Bulk Loader job. Also see configuration properties in the Configuration reference.

Read settings

Transformation settings

Function for transformScala:

Output settings

Tune performance

As the name of the Parallel Bulk Loader job implies, it is designed to ingest large amounts of data into Fusion by parallelizing the work across your Spark cluster. To achieve scalability, you might need to increase the amount of memory and/or CPU resources allocated to the job.For Spark configuration information, see Spark operations.You can pass these properties in the job configuration to override the default Spark shell options:Spark shell options

Examples

Here we provide screenshots and example JSON job definitions to illustrate key points about how to load from different data sources.

Use NLP during indexing

In this example, we leverage the John Snow labs NLP library during indexing. This is just quick-and-dirty to show the concept.Also see:Use this transform Scala script:
Be sure to add the JohnSnowLabs:spark-nlp:1.4.2 package using Spark Shell Options.

Clean up data with SQL transformations

Fusion has a Local Filesystem connector that can handle files such as CSV and JSON files. Using the Parallel Bulk Loader lets you leverage features that are not in the Local Filesystem connector, such as using SQL to clean up the input data.SQL transformation of CSV dataUse the following SQL to clean up the input data before indexing:

Read from S3

It is easy to read from an S3 bucket without pulling data down to your local workstation first. To avoid exposing your AWS credentials, add them to a file named core-site.xml in the apps/spark-dist/conf directory, such as:
Then you can load files using the S3a protocol, such as: s3a://sstk-dev/data/u.user. If you are running a Fusion cluster then each instance of Fusion will need a core-site.xml file. S3a is the preferred protocol for reading data into Spark because it uses Amazon’s libraries to read from S3 instead of the legacy Hadoop libraries. If you need other S3 protocols (for example, s3 or s3n) you will need to add the equivalent properties to core-site.xml.S3 Read OptionsYou will need to add the org.apache.hadoop:hadoop-aws:2.7.3 package to the job using the --packages Spark option. Also, you will need to exclude the com.fasterxml.jackson.core:jackson-core,joda-time:joda-time packages using the --exclude-packages option.S3 Spark Shell OptionsYou can also read from Google Cloud Storage (GCS), but you will need a few more properties in your core-site.xml; see Installing the Cloud Storage connector.

Read from Parquet

Reading from parquet files is built into Spark using the “parquet” format. For additional read options, see Configuration of Parquet.Read signals from parquet fileJob JSON:
This example also uses the transformScala option to filter and transform the input DataFrame into a better form for indexing using the following Scala script:

Read from JDBC tables

You can use the Parallel Bulk Loader to parallelize reads from JDBC tables, if the tables have numeric columns that can be partitioned into relatively equal partition sizes. In the example below, we partition the employees table into 4 partitions using the emp_no column (int). Behind the scenes, Spark sends four separate queries to the database and processes the result sets in parallel.

Load the JDBC driver JAR file into the Blob store

Before you ingest from a JDBC data source, you need to use the Fusion Admin UI to upload the JDBC driver JAR file into the blob store.Alternatively, you can add the JAR file to the Fusion blob store with resourceType=spark:jar; for example:
At runtime, Fusion’s Spark job management framework knows how to add any JAR files with resourceType=spark:jar from the blob store to the appropriate classpaths before running a Parallel Bulk Loader job.

Read from a table

Read from JDBC tablesFor more information on reading from JDBC-compliant databases, see:
Notice the use of the $MIN(emp_no) and $MAX(emp_no) macros in the read options. These are macros offered by the Parallel Bulk Loader to help configure parallel reads of JDBC tables. Behind the scenes, the macros are translated into SQL queries to get the MAX and MIN values of the specified field, which Spark uses to compute splits for partitioned queries. As mentioned above, the field must be numeric and must have a relatively balanced distribution of values between MAX and MIN; otherwise, you are unlikely to see much performance benefit to partitioning.

Index HBase tables

To index an HBase table, use the Hortonworks connector.
The Parallel Bulk Loader lets us replace the HBase Indexer.
You will need to add an hbase-site.xml (and possibly core-site.xml) to apps/spark-dist/conf in Fusion, for example:
For this example, we will create a test table in HBase. If you already have a table in HBase, feel free to use that table instead.
  1. Launch the HBase shell and create a table named fusion_nums with a single column family named lw:
  2. Do a list command to see your table:
  3. Fill the table with some data:
  4. Scan the fusion_nums table to see your data:
The HBase connector requires a catalog read option that defines the columns you want to read and how to map them into a Spark DataFrame. For our sample table, the following suffices:
Index HBase tablesNotice the use of the $lastTimestamp macro in the read options. This lets us filter rows read from HBase using the timestamp of the last document the Parallel Bulk Loader wrote to Solr, that is, to get the newest updates from HBase only (incremental updates). Most Spark data sources provide a way to filter results based on timestamp.Job JSON:

Index Elastic data

With Elasticsearch 6.2.2 using the org.elasticsearch:elasticsearch-spark-20_2.11:6.2.1 package, here is a Scala script to run in bin/spark-shell to index some test data:
ElasticJob JSON:

Read from Couchbase

To index a Couchbase bucket, use the official Couchbase Spark connector found here.For example, we will create a test bucket in Couchbase. If you already have a bucket in Couchbase, feel free to use that and skip to the test data setup section. This test was performed using Couchbase Server 6.0.0.
  1. Create a bucket test in the Couchbase admin UI. Give access to a system account user to use in the Parallel Bulk Loader job config.
  2. Connect to Couchbase using the command line client cbq. For example, cbq -e=http://<host>:8091 -u <user> -p <password>. Ensure the provided user is an authorized user of the test bucket.
  3. Create a primary index on the test bucket: CREATE PRIMARY INDEX 'test-primary-index' ON 'test' USING GSI;.
  4. Insert some data: INSERT INTO 'test' ( KEY, VALUE ) VALUES ( "1", { "id": "01", "field1": "a value", "field2": "another value"} ) RETURNING META().id as docid, *;.
  5. Verify you can query the document just inserted: select * from 'test';.
To ingest from this bucket with the Parallel Bulk Loader, use the Couchbase Spark connector by specifying the format com.couchbase.spark.sql.DefaultSource. Then specify the com.couchbase.client:spark-connector_2.11:2.2.0 package as the spark shell --packages option, as well as a few spark settings that direct the connector to a particular Couchbase server and bucket to connect to using the provided credentials. See here for all of the available Spark configuration settings for the Couchbase Spark connector.Putting it all together:Couchbase

XML setup

XML is a supported format that requires settings for format and --packages. In addition, you must specify the filepath in the readOptions section. For example:

API usage examples

The following examples include requests and responses for some of the API endpoints.

Create a parallel bulk loader job

An example request to create a parallel bulk loader job with the REST API is as follows:
The response is a message indicating success or failure.

Start a parallel bulk loader job

This request starts the pbl_load PBL job.
The response is:

POST spark configuration to a parallel bulk loader job

This request posts Spark configuration information to the store_typeahead_entity_load PBL job.

Learn more

You can use the Parallel Bulk Loader to copy Solr collections from one collection to another. This is helpful for copying collections from a production environment to a testing or development environment and using real data in the development and testing process.

Parallel Bulk Loader Job Configuration

In the Lucidworks Search UI:
  1. Navigate to Collections > Jobs.
  2. Click Add and select Parallel Bulk Loader from the menu.
  3. Enter a name for your job, and enter solr as the format.
  4. Set the parameter name and value for the collection you want to import.
    • Parameter name: collection
    • Parameter value: the name of the collection to import from
  5. If the source collection is in a different Lucidworks Search app, add an additional parameter name and value pair.
    • Parameter name: zkHost
    • Parameter value: the location of zkHost
  6. Enter the output collection, or where you want to import the collection to.
    • Output collection: the name of the output collection
    • Send to index pipeline: _system
  7. Save your job.
When you’re ready, run the job or schedule the job to run at an interval of your choice.

Parallel Bulk Loader JSON Configuration

Use this sample JSON payload to load the Parallel Bulk Loader to the Spark Jobs API:
Replace SOURCE_COLLECTION with the name of the collection being exported. Replace OUTPUT_COLLECTION with the name of the destination collection.Use the following Spark Jobs API call to run the job. Replace data.json with your JSON payload file.

Create and run Parallel Bulk Loader jobs

Use the Jobs manager to create and run Parallel Bulk Loader jobs. You can also use the Scheduler to schedule jobs.In the procedures, select Parallel Bulk Loader as the job type and configure the job using the following information.

Configuration settings for the Parallel Bulk Loader job

This section provides configuration settings for the Parallel Bulk Loader job.

Read settings

Transformation settings

Function for Parallel Bulk Loader:

Output settings

Tune performance

As the name of the Parallel Bulk Loader job implies, it is designed to ingest large amounts of data into Fusion by parallelizing the work across your Spark cluster. To achieve scalability, you might need to increase the amount of memory and/or CPU resources allocated to the job.You can pass these properties in the job configuration to override the default Spark shell options:Spark shell options

Examples

Here we provide screenshots and example JSON job definitions to illustrate key points about how to load from different data sources.

Use NLP during indexing

In this example, we leverage the John Snow labs NLP library during indexing. This is just quick-and-dirty to show the concept.Also see:Use this transform Scala script:
Be sure to add the JohnSnowLabs:spark-nlp:1.4.2 package using Spark Shell Options.

Clean up data with SQL transformations

Fusion has a Local Filesystem connector that can handle files such as CSV and JSON files. Using the Parallel Bulk Loader lets you leverage features that are not in the Local Filesystem connector, such as using SQL to clean up the input data.SQL transformation of CSV dataUse the following SQL to clean up the input data before indexing:

Read from S3

It is easy to read from an S3 bucket without pulling data down to your local workstation first. To avoid exposing your AWS credentials, add them to a file named core-site.xml in the apps/spark-dist/conf directory, such as:
Then you can load files using the S3a protocol, such as: s3a://sstk-dev/data/u.user. If you are running a Fusion cluster then each instance of Fusion will need a core-site.xml file. S3a is the preferred protocol for reading data into Spark because it uses Amazon’s libraries to read from S3 instead of the legacy Hadoop libraries. If you need other S3 protocols (for example, s3 or s3n) you will need to add the equivalent properties to core-site.xml.S3 Read OptionsYou will need to add the org.apache.hadoop:hadoop-aws:2.7.3 package to the job using the --packages Spark option. Also, you will need to exclude the com.fasterxml.jackson.core:jackson-core,joda-time:joda-time packages using the --exclude-packages option.S3 Spark Shell OptionsYou can also read from Google Cloud Storage (GCS), but you will need a few more properties in your core-site.xml; see Installing the Cloud Storage connector.

Read from Parquet

Reading from parquet files is built into Spark using the “parquet” format. For additional read options, see Configuration of Parquet.Read signals from parquet fileJob JSON:
This example also uses the transformScala option to filter and transform the input DataFrame into a better form for indexing using the following Scala script:

Read from JDBC tables

You can use the Parallel Bulk Loader to parallelize reads from JDBC tables, if the tables have numeric columns that can be partitioned into relatively equal partition sizes. In the example below, we partition the employees table into 4 partitions using the emp_no column (int). Behind the scenes, Spark sends four separate queries to the database and processes the result sets in parallel.

Load the JDBC driver JAR file into the Blob store

Before you ingest from a JDBC data source, you need to use the Fusion Admin UI to upload the JDBC driver JAR file into the blob store.Alternatively, you can add the JAR file to the Fusion blob store with resourceType=spark:jar; for example:
At runtime, Fusion’s Spark job management framework knows how to add any JAR files with resourceType=spark:jar from the blob store to the appropriate classpaths before running a Parallel Bulk Loader job.

Read from a table

Read from JDBC tablesFor more information on reading from JDBC-compliant databases, see:
Notice the use of the $MIN(emp_no) and $MAX(emp_no) macros in the read options. These are macros offered by the Parallel Bulk Loader to help configure parallel reads of JDBC tables. Behind the scenes, the macros are translated into SQL queries to get the MAX and MIN values of the specified field, which Spark uses to compute splits for partitioned queries. As mentioned above, the field must be numeric and must have a relatively balanced distribution of values between MAX and MIN; otherwise, you are unlikely to see much performance benefit to partitioning.

Index HBase tables

To index an HBase table, use the Hortonworks connector.
The Parallel Bulk Loader lets us replace the HBase Indexer.
You will need to add an hbase-site.xml (and possibly core-site.xml) to apps/spark-dist/conf in Fusion, for example:
For this example, we will create a test table in HBase. If you already have a table in HBase, feel free to use that table instead.
  1. Launch the HBase shell and create a table named fusion_nums with a single column family named lw:
  2. Do a list command to see your table:
  3. Fill the table with some data:
  4. Scan the fusion_nums table to see your data:
The HBase connector requires a catalog read option that defines the columns you want to read and how to map them into a Spark DataFrame. For our sample table, the following suffices:
Index HBase tablesNotice the use of the $lastTimestamp macro in the read options. This lets us filter rows read from HBase using the timestamp of the last document the Parallel Bulk Loader wrote to Solr, that is, to get the newest updates from HBase only (incremental updates). Most Spark data sources provide a way to filter results based on timestamp.Job JSON:

Index Elastic data

With Elasticsearch 6.2.2 using the org.elasticsearch:elasticsearch-spark-20_2.11:6.2.1 package, here is a Scala script to run in bin/spark-shell to index some test data:
ElasticJob JSON:

Read from Couchbase

To index a Couchbase bucket, use the official Couchbase Spark connector found here.For example, we will create a test bucket in Couchbase. If you already have a bucket in Couchbase, feel free to use that and skip to the test data setup section. This test was performed using Couchbase Server 6.0.0.
  1. Create a bucket test in the Couchbase admin UI. Give access to a system account user to use in the Parallel Bulk Loader job config.
  2. Connect to Couchbase using the command line client cbq. For example, cbq -e=http://<host>:8091 -u <user> -p <password>. Ensure the provided user is an authorized user of the test bucket.
  3. Create a primary index on the test bucket: CREATE PRIMARY INDEX 'test-primary-index' ON 'test' USING GSI;.
  4. Insert some data: INSERT INTO 'test' ( KEY, VALUE ) VALUES ( "1", { "id": "01", "field1": "a value", "field2": "another value"} ) RETURNING META().id as docid, *;.
  5. Verify you can query the document just inserted: select * from 'test';.
To ingest from this bucket with the Parallel Bulk Loader, use the Couchbase Spark connector by specifying the format com.couchbase.spark.sql.DefaultSource. Then specify the com.couchbase.client:spark-connector_2.11:2.2.0 package as the spark shell --packages option, as well as a few spark settings that direct the connector to a particular Couchbase server and bucket to connect to using the provided credentials. See here for all of the available Spark configuration settings for the Couchbase Spark connector.Putting it all together:Couchbase

XML setup

XML is a supported format that requires settings for format and --packages. In addition, you must specify the filepath in the readOptions section. For example:

Configuration properties